Class: Fluent::GcloudPubSubInput

Inherits:
Input
  • Object
show all
Defined in:
lib/fluent/plugin/in_gcloud_pubsub.rb

Defined Under Namespace

Classes: RPCServlet

Instance Method Summary collapse

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