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
79 80 81 82 83 84 85 86 87 |
# File 'lib/fluent/plugin/in_gcloud_pubsub.rb', line 79 def configure(conf) super @rpc_srv = nil @rpc_thread = nil @stop_pull = false @parser = Plugin.new_parser(@format) @parser.configure(conf) end |
#shutdown ⇒ Object
100 101 102 103 104 105 106 107 108 109 110 111 112 |
# File 'lib/fluent/plugin/in_gcloud_pubsub.rb', line 100 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_thread.join end |
#start ⇒ Object
89 90 91 92 93 94 95 96 97 98 |
# File 'lib/fluent/plugin/in_gcloud_pubsub.rb', line 89 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_thread = Thread.new(&method(:subscribe)) end |
#start_pull ⇒ Object
119 120 121 122 |
# File 'lib/fluent/plugin/in_gcloud_pubsub.rb', line 119 def start_pull @stop_pull = false log.info "start pull from subscription:#{@subscription}" end |
#stop_pull ⇒ Object
114 115 116 117 |
# File 'lib/fluent/plugin/in_gcloud_pubsub.rb', line 114 def stop_pull @stop_pull = true log.info "stop pull from subscription:#{@subscription}" end |