Class: Fluent::Plugin::SQSInput
- Inherits:
-
Input
- Object
- Input
- Fluent::Plugin::SQSInput
- Defined in:
- lib/fluent/plugin/in_sqs.rb
Instance Method Summary collapse
- #client ⇒ Object
- #configure(conf) ⇒ Object
- #parse_message(message) ⇒ Object
- #prase_json_string(jsonStr) ⇒ Object
- #queue ⇒ Object
- #run ⇒ Object
- #shutdown ⇒ Object
- #start ⇒ Object
Instance Method Details
#client ⇒ Object
52 53 54 |
# File 'lib/fluent/plugin/in_sqs.rb', line 52 def client @client ||= Aws::SQS::Client.new(stub_responses: @stub_responses) end |
#configure(conf) ⇒ Object
28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 |
# File 'lib/fluent/plugin/in_sqs.rb', line 28 def configure(conf) super if @tag == nil raise Fluent::ConfigError, "tag configuration key is mandatory" end if @sqs_url == nil raise Fluent::ConfigError, "sqs_url configuration key is mandatory" end Aws.config = { access_key_id: @aws_key_id, secret_access_key: @aws_sec_key, region: @region } end |
#parse_message(message) ⇒ Object
95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 |
# File 'lib/fluent/plugin/in_sqs.rb', line 95 def () record = {} if @parse_body_as_json == true record['body'] = prase_json_string(.body.to_s) else record['body'] = .body.to_s end if @add_receipt_handle == true record['receipt_handle'] = .receipt_handle.to_s end if @add_message_id == true record['message_id'] = ..to_s end if @add_md5_of_body == true record['md5_of_body'] = .md5_of_body.to_s end if @add_queue_url == true record['queue_url'] = .queue_url.to_s end if @attribute_name_to_extract.to_s.strip.length > 0 record[@attribute_name_to_extract] = .[@attribute_name_to_extract].string_value.to_s end record end |
#prase_json_string(jsonStr) ⇒ Object
83 84 85 86 87 88 89 90 91 92 93 |
# File 'lib/fluent/plugin/in_sqs.rb', line 83 def prase_json_string(jsonStr) response = jsonStr begin response = JSON.parse(jsonStr) rescue => ex # unable to pase json str (str is not a valid json) end response end |
#queue ⇒ Object
56 57 58 |
# File 'lib/fluent/plugin/in_sqs.rb', line 56 def queue @queue ||= Aws::SQS::Resource.new(client: client).queue(@sqs_url) end |
#run ⇒ Object
64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 |
# File 'lib/fluent/plugin/in_sqs.rb', line 64 def run queue.( max_number_of_messages: @max_number_of_messages, wait_time_seconds: @wait_time_seconds, visibility_timeout: @visibility_timeout, attribute_names: ['All'], # Receive all available built-in message attributes. message_attribute_names: ['All'] # Receive any custom message attributes. ).each do || record = () .delete if @delete_message router.emit(@tag, Fluent::Engine.now, record) end rescue => ex log.error 'failed to emit or receive', error: ex.to_s, error_class: ex.class log.warn_backtrace ex.backtrace end |
#shutdown ⇒ Object
60 61 62 |
# File 'lib/fluent/plugin/in_sqs.rb', line 60 def shutdown super end |
#start ⇒ Object
46 47 48 49 50 |
# File 'lib/fluent/plugin/in_sqs.rb', line 46 def start super timer_execute(:in_sqs_run_periodic_timer, @receive_interval, &method(:run)) end |