Class: Fluent::Plugin::NetflowipfixInput::ParserThread

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

Overview

class UdpListenerThread

Instance Method Summary collapse

Constructor Details

#initialize(udpQueue, queuesleep, eventQueue, cache_ttl, definitions, log) ⇒ ParserThread

Returns a new instance of ParserThread.



247
248
249
250
251
252
253
254
255
256
257
258
259
# File 'lib/fluent/plugin/in_netflowipfix.rb', line 247

def initialize(udpQueue, queuesleep, eventQueue, cache_ttl, definitions, log)
  @udpQueue = udpQueue
  @queuesleep = queuesleep
  @eventQueue = eventQueue
  @log = log

  @parser_v5 = NetflowipfixInput::ParserNetflowv5.new
  @parser_v9 = NetflowipfixInput::ParserNetflowv9.new
  @parser_v10 = NetflowipfixInput::ParserIPfixv10.new

  @parser_v9.configure(cache_ttl, definitions)
  @parser_v10.configure(cache_ttl, definitions)
end

Instance Method Details

#closeObject



265
266
267
268
269
270
271
# File 'lib/fluent/plugin/in_netflowipfix.rb', line 265

def close
  # Garbage collection
  @parser_v5 = nil
  @parser_v9 = nil
  @parser_v10 = nil
  GC.start
end

#emit(time, event, host = nil) ⇒ Object

def run



318
319
320
321
322
323
324
# File 'lib/fluent/plugin/in_netflowipfix.rb', line 318

def emit(time, event, host = nil)
  if !host.nil?
    event["host"] = host
  end
  @eventQueue << [time, event]
    @log.trace "ParserThread::emit #{@eventQueue.length}"
end

#joinObject



273
274
275
# File 'lib/fluent/plugin/in_netflowipfix.rb', line 273

def join
    @thread.join
end

#runObject



277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
# File 'lib/fluent/plugin/in_netflowipfix.rb', line 277

def run
  loop do
    if @udpQueue.length == 0
      sleep(@queuesleep)

    else
      block = method(:emit)
      ar = @udpQueue.pop
      time = ar[0]
      msg = ar[1]
      payload = msg["message"]
      host = msg["sender"]
      
      version,_ = payload[0,2].unpack('n')
      @log.trace "ParserThread::pop #{@udpQueue.length} v#{version}"

      case version
        when 5          
          packet = NetflowipfixInput::Netflow5Packet.read(payload)
          @parser_v5.handle_v5(host, packet, block)
        when 9
          packet = NetflowipfixInput::Netflow9Packet.read(payload)
          @parser_v9.handle_v9(host, packet, block)
        when 10
          packet = NetflowipfixInput::Netflow10Packet.read(payload)
          @parser_v10.handle_v10(host, packet, block)
        else
          $log.warn "Unsupported Netflow version v#{version}: #{version.class}"
      end # case

      # Free up variables for garbage collection
      ar = @udpQueue.pop
      version = nil
      time = nil
      msg = nil
      payload = nil
      host = nil

    end
  end # loop do
end

#startObject



260
261
262
263
# File 'lib/fluent/plugin/in_netflowipfix.rb', line 260

def start
  @thread = Thread.new(&method(:run))
  @log.debug "ParserThread::start"
end