Class: Kafka::AsyncProducer::Worker
- Inherits:
-
Object
- Object
- Kafka::AsyncProducer::Worker
- Defined in:
- lib/kafka/async_producer.rb
Instance Method Summary collapse
-
#initialize(queue:, producer:, delivery_threshold:) ⇒ Worker
constructor
A new instance of Worker.
- #run ⇒ Object
Constructor Details
#initialize(queue:, producer:, delivery_threshold:) ⇒ Worker
Returns a new instance of Worker.
183 184 185 186 187 |
# File 'lib/kafka/async_producer.rb', line 183 def initialize(queue:, producer:, delivery_threshold:) @queue = queue @producer = producer @delivery_threshold = delivery_threshold end |
Instance Method Details
#run ⇒ Object
189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 |
# File 'lib/kafka/async_producer.rb', line 189 def run loop do operation, payload = @queue.pop case operation when :produce produce(*payload) if threshold_reached? when :deliver_messages when :shutdown # Deliver any pending messages first. # Stop the run loop. break else raise "Unknown operation #{operation.inspect}" end end ensure @producer.shutdown end |