Class: Fluent::GcloudPubSubInput
- Inherits:
-
Input
- Object
- Input
- Fluent::GcloudPubSubInput
show all
- Defined in:
- lib/fluent/plugin/in_gcloud_pubsub.rb
Defined Under Namespace
Classes: FailedParseError, RPCServlet
Instance Method Summary
collapse
Instance Method Details
103
104
105
106
107
108
109
110
111
|
# File 'lib/fluent/plugin/in_gcloud_pubsub.rb', line 103
def configure(conf)
super
@rpc_srv = nil
@rpc_thread = nil
@stop_pull = false
@parser = Plugin.new_parser(@format)
@parser.configure(conf)
end
|
#desc(description) ⇒ Object
15
16
|
# File 'lib/fluent/plugin/in_gcloud_pubsub.rb', line 15
def desc(description)
end
|
#shutdown ⇒ Object
127
128
129
130
131
132
133
134
135
136
137
138
139
|
# File 'lib/fluent/plugin/in_gcloud_pubsub.rb', line 127
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
113
114
115
116
117
118
119
120
121
122
123
124
125
|
# File 'lib/fluent/plugin/in_gcloud_pubsub.rb', line 113
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
146
147
148
149
|
# File 'lib/fluent/plugin/in_gcloud_pubsub.rb', line 146
def start_pull
@stop_pull = false
log.info "start pull from subscription:#{@subscription}"
end
|
#stop_pull ⇒ Object
141
142
143
144
|
# File 'lib/fluent/plugin/in_gcloud_pubsub.rb', line 141
def stop_pull
@stop_pull = true
log.info "stop pull from subscription:#{@subscription}"
end
|