Module: KinesisShard
Defined Under Namespace
Classes: MemoryStateStore, StateStore
Instance Method Summary
collapse
Instance Method Details
#emit_records(data, shard_id) ⇒ Object
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
|
# File 'lib/fluent/plugin/kinesis_shard.rb', line 60
def emit_records(data, shard_id)
me = Fluent::MultiEventStream.new
data.each do |d|
if @use_base64
d = Base64.decode64(d)
end
time, record = @parser.parse(d)
if record.nil? || record.empty?
$log.warn "format error :=> record #{time} : #{d}"
else
me.add(time, record)
end
end
unless me.empty?
router.emit_stream(@tag, me)
end
end
|
#get_shard_iterator_info(shard_id = '', last_sequence_number = '') ⇒ Object
46
47
48
49
50
51
52
53
54
55
56
57
58
|
# File 'lib/fluent/plugin/kinesis_shard.rb', line 46
def get_shard_iterator_info(shard_id='', last_sequence_number='')
if last_sequence_number.empty?
shard_iterator_info = @client.get_shard_iterator(
stream_name: @stream_name, shard_id: shard_id, shard_iterator_type: 'TRIM_HORIZON')
else
shard_iterator_info = @client.get_shard_iterator(
stream_name: @stream_name, shard_id: shard_id, shard_iterator_type: 'AFTER_SEQUENCE_NUMBER', starting_sequence_number: last_sequence_number)
end
rescue => e
$log.warn "does not AFTER_SEQUENCE_NUMBER : #{e.message}"
shard_iterator_info = @client.get_shard_iterator(
stream_name: @stream_name, shard_id: shard_id, shard_iterator_type: 'TRIM_HORIZON')
end
|
#load_records_thread(shard_id) ⇒ Object
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
|
# File 'lib/fluent/plugin/kinesis_shard.rb', line 5
def load_records_thread(shard_id)
begin
state_store = @state_dir_path.nil? ? MemoryStateStore.new : StateStore.new(@state_dir_path, shard_id)
rescue => e
$log.warn "does not StateStore !!: #{e.message}"
state_store = MemoryStateStore.new
end
last_sequence_number = state_store.load_sequence_number
shard_iterator_info = get_shard_iterator_info(shard_id, last_sequence_number)
shard_iterator = shard_iterator_info.shard_iterator
while !@stop_flag && !@thread_stop_map[shard_id] do
begin
records_info = @client.get_records(shard_iterator: shard_iterator, limit: @load_records_limit)
rescue => e
$log.error "get record Error: #{e.message}"
re_shard_iterator_info = get_shard_iterator_info(shard_id, last_sequence_number)
records_info = @client.get_records(
shard_iterator: re_shard_iterator_info.shard_iterator, limit: @load_records_limit/10)
end
if records_info.next_shard_iterator.nil?
@thread_stop_map[shard_id] = true
break
end
data = records_info.records.map(&:data)
emit_records(data, shard_id)
tmp_last_sequence_number = sequence(records_info)
unless tmp_last_sequence_number.nil?
state_store.update(tmp_last_sequence_number)
last_sequence_number = tmp_last_sequence_number
end
shard_iterator = records_info.next_shard_iterator
sleep(@load_record_interval)
end
end
|
#sequence(records_info) ⇒ Object
80
81
82
83
84
85
86
87
|
# File 'lib/fluent/plugin/kinesis_shard.rb', line 80
def sequence(records_info)
sequence_number_list = records_info.records.map(&:sequence_number)
if sequence_number_list.length > 0
sequence_number = records_info.records.map(&:sequence_number)[-1]
else
sequence_number = nil
end
end
|