Class: Kcl::Proxies::KinesisProxy

Inherits:
Object
  • Object
show all
Defined in:
lib/kcl/proxies/kinesis_proxy.rb

Instance Attribute Summary collapse

Instance Method Summary collapse

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

#clientObject (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

Parameters:

  • shard_iterator (String)

Returns:

  • (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

Parameters:

  • shard_id (String)
  • shard_iterator_type (String) (defaults to: nil)

Returns:

  • (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

Parameters:

  • data (Hash)

Returns:

  • (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

#shardsArray

Returns:

  • (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