Class: Fluent::GcloudPubSubInput
- Inherits:
-
Input
- Object
- Input
- Fluent::GcloudPubSubInput
- Defined in:
- lib/fluent/plugin/in_gcloud_pubsub.rb
Defined Under Namespace
Classes: RPCServlet
Instance Method Summary collapse
- #configure(conf) ⇒ Object
- #shutdown ⇒ Object
- #start ⇒ Object
- #start_pull ⇒ Object
- #stop_pull ⇒ Object
Instance Method Details
#configure(conf) ⇒ Object
80 81 82 83 84 85 86 87 88 |
# File 'lib/fluent/plugin/in_gcloud_pubsub.rb', line 80 def configure(conf) super @rpc_srv = nil @rpc_thread = nil @stop_pull = false @parser = Plugin.new_parser(@format) @parser.configure(conf) end |
#shutdown ⇒ Object
104 105 106 107 108 109 110 111 112 113 114 115 116 |
# File 'lib/fluent/plugin/in_gcloud_pubsub.rb', line 104 def shutdown super if @rpc_srv @rpc_srv.shutdown @rpc_srv = nil end if @rpc_thread @rpc_thread.join @rpc_thread = nil end @stop_subscribing = true @subscribe_threads.each(&:join) end |
#start ⇒ Object
90 91 92 93 94 95 96 97 98 99 100 101 102 |
# File 'lib/fluent/plugin/in_gcloud_pubsub.rb', line 90 def start super start_rpc if @enable_rpc @subscriber = Fluent::GcloudPubSub::Subscriber.new @project, @key, @topic, @subscription log.debug "connected subscription:#{@subscription} in project #{@project}" @stop_subscribing = false @subscribe_threads = [] @pull_threads.times do @subscribe_threads.push Thread.new(&method(:subscribe)) end end |
#start_pull ⇒ Object
123 124 125 126 |
# File 'lib/fluent/plugin/in_gcloud_pubsub.rb', line 123 def start_pull @stop_pull = false log.info "start pull from subscription:#{@subscription}" end |
#stop_pull ⇒ Object
118 119 120 121 |
# File 'lib/fluent/plugin/in_gcloud_pubsub.rb', line 118 def stop_pull @stop_pull = true log.info "stop pull from subscription:#{@subscription}" end |