Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion fluent-plugin-sqs.gemspec
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@

Gem::Specification.new do |s|
s.name = 'fluent-plugin-sqs'
s.version = '3.0.0'
s.version = '3.0.1'

s.required_rubygems_version = Gem::Requirement.new('>= 0') if s.respond_to? :required_rubygems_version=
s.authors = ['Yuri Odagiri']
Expand Down
19 changes: 17 additions & 2 deletions lib/fluent/plugin/in_sqs.rb
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
require 'fluent/plugin/input'
require 'aws-sdk-sqs'
require 'json'

module Fluent::Plugin
class SQSInput < Input
Expand All @@ -10,6 +11,7 @@ class SQSInput < Input
config_param :aws_key_id, :string, default: nil, secret: true
config_param :aws_sec_key, :string, default: nil, secret: true
config_param :tag, :string
config_param :tag_key, :string, default: nil
config_param :region, :string, default: 'ap-northeast-1'
config_param :sqs_url, :string, default: nil
config_param :receive_interval, :time, default: 0.1
Expand All @@ -18,6 +20,7 @@ class SQSInput < Input
config_param :visibility_timeout, :integer, default: nil
config_param :delete_message, :bool, default: false
config_param :stub_responses, :bool, default: false
config_param :raw_message, :bool, default: false

def configure(conf)
super
Expand Down Expand Up @@ -53,11 +56,15 @@ def run
wait_time_seconds: @wait_time_seconds,
visibility_timeout: @visibility_timeout
).each do |message|
record = parse_message(message)
record = @raw_message ? parse_raw_message(message) : parse_message(message)

message.delete if @delete_message

router.emit(@tag, Fluent::Engine.now, record)
tag = @tag_key ? record[@tag_key] : @tag
record.delete @tag_key if @tag_key

log.debug "emiting record under tag #{tag}: #{record.to_s}"
router.emit(tag, Fluent::Engine.now, record)
end
rescue
log.error 'failed to emit or receive', error: $ERROR_INFO.to_s, error_class: $ERROR_INFO.class.to_s
Expand All @@ -66,6 +73,14 @@ def run

private

def get_tag_name(record)
record[@tag_key] rescue @tag
end

def parse_raw_message(message)
JSON.parse(message.body) rescue message.body.to_s
end

def parse_message(message)
{
'body' => message.body.to_s,
Expand Down