You are viewing documentation for version 3 of the AWS SDK for Ruby. Version 2 documentation can be found here.

Class: Aws::Kinesis::AsyncClient

Inherits:
Seahorse::Client::AsyncBase show all
Includes:
AsyncClientStubs
Defined in:
gems/aws-sdk-kinesis/lib/aws-sdk-kinesis/async_client.rb

Instance Attribute Summary

Attributes inherited from Seahorse::Client::AsyncBase

#connection

Attributes inherited from Seahorse::Client::Base

#config, #handlers

API Operations collapse

Instance Method Summary collapse

Methods included from AsyncClientStubs

#send_events

Methods included from ClientStubs

#api_requests, #stub_data, #stub_responses

Methods inherited from Seahorse::Client::AsyncBase

#close_connection, #connection_errors, #new_connection, #operation_names

Methods inherited from Seahorse::Client::Base

add_plugin, api, clear_plugins, define, new, #operation_names, plugins, remove_plugin, set_api, set_plugins

Methods included from Seahorse::Client::HandlerBuilder

#handle, #handle_request, #handle_response

Constructor Details

#initialize(*args) ⇒ AsyncClient

@option options [required, Aws::CredentialProvider] :credentials Your AWS credentials. This can be an instance of any one of the following classes:

* `Aws::Credentials` - Used for configuring static, non-refreshing
  credentials.

* `Aws::InstanceProfileCredentials` - Used for loading credentials
  from an EC2 IMDS on an EC2 instance.

* `Aws::SharedCredentials` - Used for loading credentials from a
  shared file, such as `~/.aws/config`.

* `Aws::AssumeRoleCredentials` - Used when you need to assume a role.

When `:credentials` are not configured directly, the following
locations will be searched for credentials:

* `Aws.config[:credentials]`
* The `:access_key_id`, `:secret_access_key`, and `:session_token` options.
* ENV['AWS_ACCESS_KEY_ID'], ENV['AWS_SECRET_ACCESS_KEY']
* `~/.aws/credentials`
* `~/.aws/config`
* EC2 IMDS instance profile - When used by default, the timeouts are
  very aggressive. Construct and pass an instance of
  `Aws::InstanceProfileCredentails` to enable retries and extended
  timeouts.

@option options [required, String] :region The AWS region to connect to. The configured :region is used to determine the service :endpoint. When not passed, a default :region is search for in the following locations:

* `Aws.config[:region]`
* `ENV['AWS_REGION']`
* `ENV['AMAZON_REGION']`
* `ENV['AWS_DEFAULT_REGION']`
* `~/.aws/credentials`
* `~/.aws/config`

@option options [String] :access_key_id

@option options [Boolean] :convert_params (true) When true, an attempt is made to coerce request parameters into the required types.

@option options [String] :endpoint The client endpoint is normally constructed from the :region option. You should only configure an :endpoint when connecting to test endpoints. This should be avalid HTTP(S) URI.

@option options [Proc] :event_stream_handler When an EventStream or Proc object is provided, it will be used as callback for each chunk of event stream response received along the way.

@option options [Proc] :input_event_stream_handler When an EventStream or Proc object is provided, it can be used for sending events for the event stream.

@option options [Aws::Log::Formatter] :log_formatter (Aws::Log::Formatter.default) The log formatter.

@option options [Symbol] :log_level (:info) The log level to send messages to the :logger at.

@option options [Logger] :logger The Logger instance to send log messages to. If this option is not set, logging will be disabled.

@option options [Proc] :output_event_stream_handler When an EventStream or Proc object is provided, it will be used as callback for each chunk of event stream response received along the way.

@option options [String] :profile ("default") Used when loading credentials from the shared credentials file at HOME/.aws/credentials. When not specified, 'default' is used.

@option options [Float] :retry_base_delay (0.3) The base delay in seconds used by the default backoff function.

@option options [Symbol] :retry_jitter (:none) A delay randomiser function used by the default backoff function. Some predefined functions can be referenced by name - :none, :equal, :full, otherwise a Proc that takes and returns a number.

@see https://www.awsarchitectureblog.com/2015/03/backoff.html

@option options [Integer] :retry_limit (3) The maximum number of times to retry failed requests. Only ~ 500 level server errors and certain ~ 400 level client errors are retried. Generally, these are throttling errors, data checksum errors, networking errors, timeout errors and auth errors from expired credentials.

@option options [Integer] :retry_max_delay (0) The maximum number of seconds to delay between retries (0 for no limit) used by the default backoff function.

@option options [String] :secret_access_key

@option options [String] :session_token

@option options [Boolean] :simple_json (false) Disables request parameter conversion, validation, and formatting. Also disable response data type conversions. This option is useful when you want to ensure the highest level of performance by avoiding overhead of walking request parameters and response data structures.

When `:simple_json` is enabled, the request parameters hash must
be formatted exactly as the DynamoDB API expects.

@option options [Boolean] :stub_responses (false) Causes the client to return stubbed responses. By default fake responses are generated and returned. You can specify the response data to return or errors to raise by calling ClientStubs#stub_responses. See ClientStubs for more information.

** Please note ** When response stubbing is enabled, no HTTP
requests are made, and retries are disabled.

@option options [Boolean] :validate_params (true) When true, request parameters are validated before sending the request.



179
180
181
182
183
184
# File 'gems/aws-sdk-kinesis/lib/aws-sdk-kinesis/async_client.rb', line 179

def initialize(*args)
  unless Kernel.const_defined?("HTTP2")
    raise "Must include http/2 gem to use AsyncClient instances."
  end
  super
end

Instance Method Details

#subscribe_to_shard(params = {}) ⇒ Types::SubscribeToShardOutput

Call this operation from your consumer after you call RegisterStreamConsumer to register the consumer with Kinesis Data Streams. If the call succeeds, your consumer starts receiving events of type SubscribeToShardEvent for up to 5 minutes, after which time you need to call SubscribeToShard again to renew the subscription if you want to continue to receive records.

You can make one call to SubscribeToShard per second per ConsumerARN. If your call succeeds, and then you call the operation again less than 5 seconds later, the second call generates a ResourceInUseException. If you call the operation a second time more than 5 seconds after the first call succeeds, the second call succeeds and the first connection gets shut down.

Examples:

EventStream Operation Example


You can process event once it arrives immediately, or wait until
full response complete and iterate through eventstream enumerator.

To interact with event immediately, you need to register #subscribe_to_shard
with callbacks, callbacks can be register for specifc events or for all events,
callback for errors in the event stream is also available for register.

Callbacks can be passed in by `:event_stream_handler` option or within block
statement attached to #subscribe_to_shard call directly. Hybrid pattern of both
is also supported.

`:event_stream_handler` option takes in either Proc object or
Aws::Kinesis::EventStreams::SubscribeToShardEventStream object.

Usage pattern a): callbacks with a block attached to #subscribe_to_shard
  Example for registering callbacks for all event types and error event

  client.subscribe_to_shard( # params input# ) do |stream|
    stream.on_error_event do |event|
      # catch unmodeled error event in the stream
      raise event
      # => Aws::Errors::EventError
      # event.event_type => :error
      # event.error_code => String
      # event.error_message => String
    end

    stream.on_event do |event|
      # process all events arrive
      puts event.event_type
      ...
    end

  end

Usage pattern b): pass in `:event_stream_handler` for #subscribe_to_shard

  1) create a Aws::Kinesis::EventStreams::SubscribeToShardEventStream object
  Example for registering callbacks with specific events

    handler = Aws::Kinesis::EventStreams::SubscribeToShardEventStream.new
    handler.on_subscribe_to_shard_event_event do |event|
      event # => Aws::Kinesis::Types::SubscribeToShardEvent
    end
    handler.on_resource_not_found_exception_event do |event|
      event # => Aws::Kinesis::Types::ResourceNotFoundException
    end
    handler.on_resource_in_use_exception_event do |event|
      event # => Aws::Kinesis::Types::ResourceInUseException
    end
    handler.on_kms_disabled_exception_event do |event|
      event # => Aws::Kinesis::Types::KMSDisabledException
    end
    handler.on_kms_invalid_state_exception_event do |event|
      event # => Aws::Kinesis::Types::KMSInvalidStateException
    end
    handler.on_kms_access_denied_exception_event do |event|
      event # => Aws::Kinesis::Types::KMSAccessDeniedException
    end
    handler.on_kms_not_found_exception_event do |event|
      event # => Aws::Kinesis::Types::KMSNotFoundException
    end
    handler.on_kms_opt_in_required_event do |event|
      event # => Aws::Kinesis::Types::KMSOptInRequired
    end
    handler.on_kms_throttling_exception_event do |event|
      event # => Aws::Kinesis::Types::KMSThrottlingException
    end
    handler.on_internal_failure_exception_event do |event|
      event # => Aws::Kinesis::Types::InternalFailureException
    end

  client.subscribe_to_shard( # params input #, event_stream_handler: handler)

  2) use a Ruby Proc object
  Example for registering callbacks with specific events

  handler = Proc.new do |stream|
    stream.on_subscribe_to_shard_event_event do |event|
      event # => Aws::Kinesis::Types::SubscribeToShardEvent
    end
    stream.on_resource_not_found_exception_event do |event|
      event # => Aws::Kinesis::Types::ResourceNotFoundException
    end
    stream.on_resource_in_use_exception_event do |event|
      event # => Aws::Kinesis::Types::ResourceInUseException
    end
    stream.on_kms_disabled_exception_event do |event|
      event # => Aws::Kinesis::Types::KMSDisabledException
    end
    stream.on_kms_invalid_state_exception_event do |event|
      event # => Aws::Kinesis::Types::KMSInvalidStateException
    end
    stream.on_kms_access_denied_exception_event do |event|
      event # => Aws::Kinesis::Types::KMSAccessDeniedException
    end
    stream.on_kms_not_found_exception_event do |event|
      event # => Aws::Kinesis::Types::KMSNotFoundException
    end
    stream.on_kms_opt_in_required_event do |event|
      event # => Aws::Kinesis::Types::KMSOptInRequired
    end
    stream.on_kms_throttling_exception_event do |event|
      event # => Aws::Kinesis::Types::KMSThrottlingException
    end
    stream.on_internal_failure_exception_event do |event|
      event # => Aws::Kinesis::Types::InternalFailureException
    end
  end

  client.subscribe_to_shard( # params input #, event_stream_handler: handler)

Usage pattern c): hybird pattern of a) and b)

    handler = Aws::Kinesis::EventStreams::SubscribeToShardEventStream.new
    handler.on_subscribe_to_shard_event_event do |event|
      event # => Aws::Kinesis::Types::SubscribeToShardEvent
    end
    handler.on_resource_not_found_exception_event do |event|
      event # => Aws::Kinesis::Types::ResourceNotFoundException
    end
    handler.on_resource_in_use_exception_event do |event|
      event # => Aws::Kinesis::Types::ResourceInUseException
    end
    handler.on_kms_disabled_exception_event do |event|
      event # => Aws::Kinesis::Types::KMSDisabledException
    end
    handler.on_kms_invalid_state_exception_event do |event|
      event # => Aws::Kinesis::Types::KMSInvalidStateException
    end
    handler.on_kms_access_denied_exception_event do |event|
      event # => Aws::Kinesis::Types::KMSAccessDeniedException
    end
    handler.on_kms_not_found_exception_event do |event|
      event # => Aws::Kinesis::Types::KMSNotFoundException
    end
    handler.on_kms_opt_in_required_event do |event|
      event # => Aws::Kinesis::Types::KMSOptInRequired
    end
    handler.on_kms_throttling_exception_event do |event|
      event # => Aws::Kinesis::Types::KMSThrottlingException
    end
    handler.on_internal_failure_exception_event do |event|
      event # => Aws::Kinesis::Types::InternalFailureException
    end

  client.subscribe_to_shard( # params input #, event_stream_handler: handler) do |stream|
    stream.on_error_event do |event|
      # catch unmodeled error event in the stream
      raise event
      # => Aws::Errors::EventError
      # event.event_type => :error
      # event.error_code => String
      # event.error_message => String
    end
  end

Besides above usage patterns for process events when they arrive immediately, you can also
iterate through events after response complete.

Events are available at resp.event_stream # => Enumerator
For parameter input example, please refer to following request syntax

Request syntax with placeholder values


async_resp = async_client.subscribe_to_shard({
  consumer_arn: "ConsumerARN", # required
  shard_id: "ShardId", # required
  starting_position: { # required
    type: "AT_SEQUENCE_NUMBER", # required, accepts AT_SEQUENCE_NUMBER, AFTER_SEQUENCE_NUMBER, TRIM_HORIZON, LATEST, AT_TIMESTAMP
    sequence_number: "SequenceNumber",
    timestamp: Time.now,
  },
})
# => Seahorse::Client::AsyncResponse
async_resp.wait
# => Seahorse::Client::Response
# Or use async_resp.join!

Response structure


All events are available at resp.event_stream:
resp.event_stream #=> Enumerator
resp.event_stream.event_types #=> [:subscribe_to_shard_event, :resource_not_found_exception, :resource_in_use_exception, :kms_disabled_exception, :kms_invalid_state_exception, :kms_access_denied_exception, :kms_not_found_exception, :kms_opt_in_required, :kms_throttling_exception, :internal_failure_exception]

For :subscribe_to_shard_event event available at #on_subscribe_to_shard_event_event callback and response eventstream enumerator:
event.records #=> Array
event.records[0].sequence_number #=> String
event.records[0].approximate_arrival_timestamp #=> Time
event.records[0].data #=> String
event.records[0].partition_key #=> String
event.records[0].encryption_type #=> String, one of "NONE", "KMS"
event.continuation_sequence_number #=> String
event.millis_behind_latest #=> Integer

For :resource_not_found_exception event available at #on_resource_not_found_exception_event callback and response eventstream enumerator:
event.message #=> String

For :resource_in_use_exception event available at #on_resource_in_use_exception_event callback and response eventstream enumerator:
event.message #=> String

For :kms_disabled_exception event available at #on_kms_disabled_exception_event callback and response eventstream enumerator:
event.message #=> String

For :kms_invalid_state_exception event available at #on_kms_invalid_state_exception_event callback and response eventstream enumerator:
event.message #=> String

For :kms_access_denied_exception event available at #on_kms_access_denied_exception_event callback and response eventstream enumerator:
event.message #=> String

For :kms_not_found_exception event available at #on_kms_not_found_exception_event callback and response eventstream enumerator:
event.message #=> String

For :kms_opt_in_required event available at #on_kms_opt_in_required_event callback and response eventstream enumerator:
event.message #=> String

For :kms_throttling_exception event available at #on_kms_throttling_exception_event callback and response eventstream enumerator:
event.message #=> String

For :internal_failure_exception event available at #on_internal_failure_exception_event callback and response eventstream enumerator:
event.message #=> String

Parameters:

  • params (Hash) (defaults to: {})

    ({})

Options Hash (params):

  • :consumer_arn (required, String)

    For this parameter, use the value you obtained when you called RegisterStreamConsumer.

  • :shard_id (required, String)

    The ID of the shard you want to subscribe to. To see a list of all the shards for a given stream, use ListShards.

  • :starting_position (required, Types::StartingPosition)

Yields:

  • (output_event_stream_handler)

Returns:

See Also:



444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
# File 'gems/aws-sdk-kinesis/lib/aws-sdk-kinesis/async_client.rb', line 444

def subscribe_to_shard(params = {}, options = {})
  params = params.dup
  output_event_stream_handler = _event_stream_handler(
    :output,
    params.delete(:output_event_stream_handler) || params.delete(:event_stream_handler),
    EventStreams::SubscribeToShardEventStream
  )

  yield(output_event_stream_handler) if block_given?

  req = build_request(:subscribe_to_shard, params)

  req.context[:output_event_stream_handler] = output_event_stream_handler
  req.handlers.add(Aws::Binary::DecodeHandler, priority: 95)

  req.send_request(options)
end