Class: Fluent::Plugin::ConcatFilter
- Inherits:
-
Filter
- Object
- Filter
- Fluent::Plugin::ConcatFilter
- Defined in:
- lib/fluent/plugin/filter_concat.rb
Defined Under Namespace
Classes: TimeoutError
Instance Method Summary collapse
- #configure(conf) ⇒ Object
- #filter_stream(tag, es) ⇒ Object
-
#initialize ⇒ ConcatFilter
constructor
A new instance of ConcatFilter.
- #shutdown ⇒ Object
- #start ⇒ Object
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 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 |