Class: Fluent::Plugin::ConcatFilter

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

Defined Under Namespace

Classes: TimeoutError

Instance Method Summary collapse

Constructor Details

#initialize ⇒ ConcatFilter

Returns a new instance of ConcatFilter.



31
32
33
34
35
36
# File 'lib/fluent/plugin/filter_concat.rb', line 31

def initialize
  super

  @buffer = Hash.new {|h, k| h[k] = [] }
  @timeout_map = Hash.new {|h, k| h[k] = Fluent::Engine.now }
end

Instance Method Details

#configure(conf) ⇒ Object



38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
# File 'lib/fluent/plugin/filter_concat.rb', line 38

def configure(conf)
  super

  if @n_lines && @multiline_start_regexp
    raise Fluent::ConfigError, "n_lines and multiline_start_regexp are exclusive"
  end
  if @n_lines.nil? && @multiline_start_regexp.nil?
    raise Fluent::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



73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
# File 'lib/fluent/plugin/filter_concat.rb', line 73

def filter_stream(tag, es)
  new_es = Fluent::MultiEventStream.new
  es.each do |time, record|
    if /\Afluent\.(?:trace|debug|info|warn|error|fatal)\z/ =~ tag
      new_es.add(time, record)
      next
    end
    begin
      flushed_es = process(tag, time, record)
      unless flushed_es.empty?
        flushed_es.each do |_time, new_record|
          time = _time if @use_first_timestamp
          new_es.add(time, record.merge(new_record))
        end
      end
    rescue => e
      router.emit_error_event(tag, time, record, e)
    end
  end
  new_es
end

#shutdown ⇒ Object



67
68
69
70
71
# File 'lib/fluent/plugin/filter_concat.rb', line 67

def shutdown
  @finished = true
  flush_remaining_buffer
  super
end

#start ⇒ Object



61
62
63
64
65
# File 'lib/fluent/plugin/filter_concat.rb', line 61

def start
  super
  @finished = false
  timer_execute(:filter_concat_timer, 1, &method(:on_timer))
end