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
8 changes: 7 additions & 1 deletion 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 @@ -18,6 +19,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,7 +55,7 @@ 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

Expand All @@ -66,6 +68,10 @@ def run

private

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