From 07da4901ca62da73aa5bfa2f43843c37dd86846d Mon Sep 17 00:00:00 2001 From: danielporfirio Date: Tue, 3 Jul 2018 22:04:08 -0300 Subject: [PATCH] Add one new 'raw_message' parameter Add new parameter to enable the 'raw message' that returns only the body of the message discarding the other attributes. --- fluent-plugin-sqs.gemspec | 2 +- lib/fluent/plugin/in_sqs.rb | 8 +++++++- 2 files changed, 8 insertions(+), 2 deletions(-) diff --git a/fluent-plugin-sqs.gemspec b/fluent-plugin-sqs.gemspec index bba77f8..b50e381 100644 --- a/fluent-plugin-sqs.gemspec +++ b/fluent-plugin-sqs.gemspec @@ -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'] diff --git a/lib/fluent/plugin/in_sqs.rb b/lib/fluent/plugin/in_sqs.rb index a79cb04..4679cd3 100644 --- a/lib/fluent/plugin/in_sqs.rb +++ b/lib/fluent/plugin/in_sqs.rb @@ -1,5 +1,6 @@ require 'fluent/plugin/input' require 'aws-sdk-sqs' +require 'json' module Fluent::Plugin class SQSInput < Input @@ -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 @@ -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 @@ -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,