Module: WaterDrop::Producer::Buffer
- Included in:
- WaterDrop::Producer
- Defined in:
- lib/waterdrop/producer/buffer.rb
Overview
Component for buffered operations
Instance Method Summary collapse
-
#buffer(message) ⇒ Object
Adds given message into the internal producer buffer without flushing it to Kafka.
-
#buffer_many(messages) ⇒ Object
Adds given messages into the internal producer buffer without flushing them to Kafka.
-
#flush_async ⇒ Array<Rdkafka::Producer::DeliveryHandle>
Flushes the internal buffer to Kafka in an async way.
-
#flush_sync ⇒ Array<Rdkafka::Producer::DeliveryReport>
Flushes the internal buffer to Kafka in a sync way.
Instance Method Details
#buffer(message) ⇒ Object
Adds given message into the internal producer buffer without flushing it to Kafka
20 21 22 23 24 25 26 27 28 29 30 31 |
# File 'lib/waterdrop/producer/buffer.rb', line 20 def buffer() ensure_active! = middleware.run() () @monitor.instrument( 'message.buffered', producer_id: id, message: , buffer: @messages ) { @messages << } end |
#buffer_many(messages) ⇒ Object
Adds given messages into the internal producer buffer without flushing them to Kafka
39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 |
# File 'lib/waterdrop/producer/buffer.rb', line 39 def buffer_many() ensure_active! = middleware.run_many() .each { || () } @monitor.instrument( 'messages.buffered', producer_id: id, messages: , buffer: @messages ) do .each { || @messages << } end end |
#flush_async ⇒ Array<Rdkafka::Producer::DeliveryHandle>
Flushes the internal buffer to Kafka in an async way
59 60 61 62 63 64 65 66 67 |
# File 'lib/waterdrop/producer/buffer.rb', line 59 def flush_async ensure_active! @monitor.instrument( 'buffer.flushed_async', producer_id: id, messages: @messages ) { flush(false) } end |
#flush_sync ⇒ Array<Rdkafka::Producer::DeliveryReport>
Flushes the internal buffer to Kafka in a sync way
72 73 74 75 76 77 78 79 80 |
# File 'lib/waterdrop/producer/buffer.rb', line 72 def flush_sync ensure_active! @monitor.instrument( 'buffer.flushed_sync', producer_id: id, messages: @messages ) { flush(true) } end |