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