Class: Fluent::ConcatFilter

Inherits:
Filter
  • Object
show all
Defined in:
lib/fluent/plugin/filter_concat.rb

Instance Method Summary collapse

Constructor Details

#initialize ⇒ ConcatFilter

Returns a new instance of ConcatFilter.



18
19
20
21
22
# File 'lib/fluent/plugin/filter_concat.rb', line 18

def initialize
  super

  @buffer = Hash.new{|h, k| h[k] = [] }
end

Instance Method Details

#configure(conf) ⇒ Object



24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
# File 'lib/fluent/plugin/filter_concat.rb', line 24

def configure(conf)
  super

  if @n_lines && @multiline_start_regexp
    raise ConfigError, "n_lines and multiline_start_regexp are exclusive"
  end
  if @n_lines.nil? && @multiline_start_regexp.nil?
    raise ConfigError, "Either n_lines or multiline_start_regexp is required"
  end

  @mode = nil
  case
  when @n_lines
    @mode = :line
  when @multiline_start_regexp
    @mode = :regexp
    @multiline_start_regexp = Regexp.compile(@multiline_start_regexp[1..-2])
    if @multiline_end_regexp
      @multiline_end_regexp = Regexp.compile(@multiline_end_regexp[1..-2])
    end
  end
end

#filter_stream(tag, es) ⇒ Object



52
53
54
55
56
57
58
59
60
61
62
63
# File 'lib/fluent/plugin/filter_concat.rb', line 52

def filter_stream(tag, es)
  new_es = MultiEventStream.new
  es.each do |time, record|
    begin
      new_record = process(tag, time, record)
      new_es.add(time, record.merge(new_record)) if new_record
    rescue => e
      router.emit_error_event(tag, time, record, e)
    end
  end
  new_es
end

#shutdown ⇒ Object



47
48
49
50
# File 'lib/fluent/plugin/filter_concat.rb', line 47

def shutdown
  super
  flush_all_buffer
end