Class: Kcl::Proxies::KinesisProxy
- Inherits:
-
Object
- Object
- Kcl::Proxies::KinesisProxy
- Defined in:
- lib/kcl/proxies/kinesis_proxy.rb
Instance Attribute Summary collapse
-
#client ⇒ Object
readonly
Returns the value of attribute client.
Instance Method Summary collapse
- #get_records(shard_iterator) ⇒ Hash
- #get_shard_iterator(shard_id, shard_iterator_type = nil, sequence_number = nil) ⇒ String
-
#initialize(config) ⇒ KinesisProxy
constructor
A new instance of KinesisProxy.
- #put_record(data) ⇒ Hash
- #shards ⇒ Array
Constructor Details
#initialize(config) ⇒ KinesisProxy
Returns a new instance of KinesisProxy.
7 8 9 10 11 12 13 14 15 16 17 18 |
# File 'lib/kcl/proxies/kinesis_proxy.rb', line 7 def initialize(config) @client = Aws::Kinesis::Client.new( { access_key_id: config.aws_access_key_id, secret_access_key: config.aws_secret_access_key, region: config.aws_region, endpoint: config.kinesis_endpoint, ssl_verify_peer: config.use_ssl } ) @stream_name = config.kinesis_stream_name end |
Instance Attribute Details
#client ⇒ Object (readonly)
Returns the value of attribute client.
5 6 7 |
# File 'lib/kcl/proxies/kinesis_proxy.rb', line 5 def client @client end |
Instance Method Details
#get_records(shard_iterator) ⇒ Hash
44 45 46 47 |
# File 'lib/kcl/proxies/kinesis_proxy.rb', line 44 def get_records(shard_iterator) res = @client.get_records({ shard_iterator: shard_iterator }) { records: res.records, next_shard_iterator: res.next_shard_iterator } end |
#get_shard_iterator(shard_id, shard_iterator_type = nil, sequence_number = nil) ⇒ String
29 30 31 32 33 34 35 36 37 38 39 40 |
# File 'lib/kcl/proxies/kinesis_proxy.rb', line 29 def get_shard_iterator(shard_id, shard_iterator_type = nil, sequence_number = nil) params = { stream_name: @stream_name, shard_id: shard_id, shard_iterator_type: shard_iterator_type || Kcl::Checkpoints::Sentinel::LATEST } if shard_iterator_type == Kcl::Checkpoints::Sentinel::AFTER_SEQUENCE_NUMBER params[:starting_sequence_number] = sequence_number end res = @client.get_shard_iterator(params) res.shard_iterator end |
#put_record(data) ⇒ Hash
51 52 53 54 |
# File 'lib/kcl/proxies/kinesis_proxy.rb', line 51 def put_record(data) res = @client.put_record(data) { shard_id: res.shard_id, sequence_number: res.sequence_number } end |
#shards ⇒ Array
21 22 23 24 |
# File 'lib/kcl/proxies/kinesis_proxy.rb', line 21 def shards res = @client.describe_stream({ stream_name: @stream_name }) res.stream_description.shards end |