Class: Fluent::ForkOutput
- Inherits:
-
Output
- Object
- Output
- Fluent::ForkOutput
- Defined in:
- lib/fluent/plugin/out_fork.rb
Defined Under Namespace
Classes: MaxForkSizeError
Instance Method Summary collapse
- #configure(conf) ⇒ Object
- #emit(tag, es, chain) ⇒ Object
-
#initialize ⇒ ForkOutput
constructor
A new instance of ForkOutput.
Constructor Details
#initialize ⇒ ForkOutput
Returns a new instance of ForkOutput.
15 16 17 |
# File 'lib/fluent/plugin/out_fork.rb', line 15 def initialize super end |
Instance Method Details
#configure(conf) ⇒ Object
29 30 31 32 33 34 |
# File 'lib/fluent/plugin/out_fork.rb', line 29 def configure(conf) super fallbacks = %w(skip drop log) raise Fluent::ConfigError, "max_fallback must be one of #{fallbacks.inspect}" unless fallbacks.include?(@max_fallback) end |
#emit(tag, es, chain) ⇒ Object
36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 |
# File 'lib/fluent/plugin/out_fork.rb', line 36 def emit(tag, es, chain) es.each do |time, record| org_value = record[@fork_key] if org_value.nil? log.trace "#{tag} - #{time}: skip to fork #{@fork_key}=#{org_value}" next end log.trace "#{tag} - #{time}: try to fork #{@fork_key}=#{org_value}" values = [] case @fork_value_type when 'csv' values = org_value.split(@separator) when 'array' values = org_value else values = org_value end values = values.uniq unless @no_unique if @max_size && @max_size < values.size case @max_fallback when 'skip' log.warn "#{tag} - #{time}: Skip too many forked values (max=#{@max_size}) : #{org_value}" next when 'drop' log.warn "#{tag} - #{time}: Drop too many forked values (max=#{@max_size}) : #{org_value}" values = values.take(@max_size) when 'log' log.info "#{tag} - #{time}: Too many forked values (max=#{@max_size}) : #{org_value}" end end values.reject{ |value| value.to_s == '' }.each_with_index do |value, i| log.trace "#{tag} - #{time}: reemit #{@output_key}=#{value} for #{@output_tag}" new_record = record.reject{ |k, v| k == @fork_key }.merge(@output_key => value) new_record.merge!(@index_key => i) unless @index_key.nil? router.emit(@output_tag, time, new_record) end end rescue => e log.error "#{e.}: #{e.backtrace.join(', ')}" ensure chain.next end |