Ack/Reject
Acknowledge or reject individual KubeMQ queue messages in Ruby to confirm successful processing and discard failed deliveries.
Overview
Ack and reject give you per-message control over queue delivery instead of an all-or-nothing batch outcome. When a downstream receiver polls a batch, each message stays locked on the broker — invisible to other consumers — until the consumer explicitly settles it. That's what you need when one bad record in a batch shouldn't take the rest down with it.
Settlement happens through calls on the received message: ack, which permanently removes it from the queue, reject, which returns it to the queue for redelivery, and requeue, a similar return-to-queue path. Internally the broker tracks this against a receive count, which a dead-letter policy can use to stop retrying a poison message forever.
Gotchas: an unsettled message isn't gone — it snaps back to the queue once the visibility timeout expires, so a slow consumer looks identical to a rejecting one; settle every message before that deadline, and never assume a batch is fully processed until you've called ack or reject on each one individually.
Prerequisites
- KubeMQ server running on
localhost:50000 - Ruby SDK installed (
gem install kubemq)
Code
require 'kubemq'
address = ENV.fetch('KUBEMQ_ADDRESS', 'localhost:50000')
channel = 'ruby-queues.ack-reject'
begin
client = KubeMQ::QueuesClient.new(address: address, client_id: 'qs-ack-reject-example')
puts "Connected to #{address}"
4.times do |i|
label = i.even? ? 'accept' : 'reject'
msg = KubeMQ::Queues::QueueMessage.new(channel: channel, metadata: "#{label}-#{i}", body: "data-#{i}")
client.send_queue_message(msg)
end
receiver = client.create_downstream_receiver
request = KubeMQ::Queues::QueuePollRequest.new(channel: channel, max_items: 10, wait_timeout: 5)
response = receiver.poll(request)
puts "Polled #{response.messages.size} messages"
response.messages.each do |m|
if m.metadata.start_with?('accept')
m.ack
puts "Acked: #{m.metadata}"
else
m.reject
puts "Rejected: #{m.metadata}"
end
end
rescue KubeMQ::Error => e
puts "KubeMQ error: #{e.message}"
ensure
receiver&.close
client&.close
puts 'Done'
endHow It Works
- Each received message supports
ack(acknowledge),reject(negative-acknowledge), andrequeue. - Acked messages are permanently removed from the queue.
- Nacked messages are returned to the queue for redelivery.
- Review timeouts, channel names, and client IDs before running against shared environments.
Related
Was this page helpful?