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