KubeMQ
Client SDKsRubyHow-to guidesQueues

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

main.rb
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'
end

How It Works

  • Each received message supports ack (acknowledge), reject (negative-acknowledge), and requeue.
  • 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.

Was this page helpful?

On this page