Class: Fluent::GcloudPubSubOutput
- Inherits:
-
BufferedOutput
- Object
- BufferedOutput
- Fluent::GcloudPubSubOutput
- Defined in:
- lib/fluent/plugin/out_gcloud_pubsub.rb
Instance Method Summary collapse
- #configure(conf) ⇒ Object
- #desc(description) ⇒ Object
- #format(tag, time, record) ⇒ Object
- #start ⇒ Object
- #write(chunk) ⇒ Object
Instance Method Details
#configure(conf) ⇒ Object
39 40 41 42 43 |
# File 'lib/fluent/plugin/out_gcloud_pubsub.rb', line 39 def configure(conf) super @formatter = Plugin.new_formatter(@format) @formatter.configure(conf) end |
#desc(description) ⇒ Object
11 12 |
# File 'lib/fluent/plugin/out_gcloud_pubsub.rb', line 11 def desc(description) end |
#format(tag, time, record) ⇒ Object
51 52 53 |
# File 'lib/fluent/plugin/out_gcloud_pubsub.rb', line 51 def format(tag, time, record) @formatter.format(tag, time, record).to_msgpack end |
#start ⇒ Object
45 46 47 48 49 |
# File 'lib/fluent/plugin/out_gcloud_pubsub.rb', line 45 def start super @publisher = Fluent::GcloudPubSub::Publisher.new @project, @key, @topic, @autocreate_topic log.debug "connected topic:#{@topic} in project #{@project}" end |
#write(chunk) ⇒ Object
55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 |
# File 'lib/fluent/plugin/out_gcloud_pubsub.rb', line 55 def write(chunk) = [] size = 0 chunk.msgpack_each do |msg| if .length + 1 > || size + msg.bytesize > @max_total_size publish = [] size = 0 end << msg size += msg.bytesize end if .length > 0 publish end rescue Fluent::GcloudPubSub::RetryableError => ex log.warn "Retryable error occurs. Fluentd will retry.", error_message: ex.to_s, error_class: ex.class.to_s raise ex rescue => ex log.error "unexpected error", error_message: ex.to_s, error_class: ex.class.to_s log.error_backtrace raise ex end |