KubeMQ
LearnQueuesTutorials

Ack, Nack & Requeue

Master the three message settlement options — acknowledge, negative acknowledge, and requeue to another channel.

What You Will Build

A consumer that demonstrates all three message settlement options:

  • Ack — message processed successfully, remove from queue
  • Nack — processing failed, return message to queue for redelivery
  • Requeue — redirect message to a different channel

The consumer settles each message exactly one way: ack removes it, nack returns it to the same queue, requeue moves it to another channel.

Prerequisites

Acknowledge (Ack)

Use ack when processing completes successfully. The message is permanently removed from the queue.

ack_example.go
resp, err := client.PollQueue(ctx, &kubemq.PollRequest{
    Channel:            "orders",
    MaxItems:           1,
    WaitTimeoutSeconds: 5,
})
if err != nil {
    log.Fatal(err)
}
for _, m := range resp.Messages {
    fmt.Printf("Processing: %s\n", string(m.Message.Body))
}
if err := resp.AckAll(); err != nil {
    log.Fatal(err)
}
fmt.Println("Messages acknowledged and removed")
ack_example.py
response = client.receive_queue_messages(
    channel="orders",
    max_messages=1,
    wait_timeout_in_seconds=5,
)
for msg in response.messages:
    print(f"Processing: {msg.body.decode('utf-8')}")
    msg.ack()
    print("Message acknowledged and removed")
ack_example.ts
const messages = await client.receiveQueueMessages({
  channel: 'orders',
  maxMessages: 1,
  waitTimeoutSeconds: 5,
});
for (const msg of messages) {
  console.log('Processing:', new TextDecoder().decode(msg.body));
  await msg.ack();
  console.log('Message acknowledged and removed');
}
AckExample.java
ReceiveQueueMessagesResponse response = client.receiveQueueMessages(
    ReceiveQueueMessagesRequest.builder()
        .channel("orders")
        .maxMessages(1)
        .waitTimeoutSeconds(5)
        .build());

for (QueueMessageReceived msg : response.getMessages()) {
    System.out.println("Processing: " + new String(msg.getBody()));
    msg.ack();
    System.out.println("Message acknowledged and removed");
}
AckExample.cs
var response = await client.ReceiveQueueMessagesAsync(new ReceiveQueueMessagesRequest
{
    Channel = "orders",
    MaxMessages = 1,
    WaitTimeoutSeconds = 5,
});
foreach (var msg in response.Messages)
{
    Console.WriteLine($"Processing: {Encoding.UTF8.GetString(msg.Body.Span)}");
    await msg.AckAsync();
    Console.WriteLine("Message acknowledged and removed");
}
AckExample.kt
val response = client.receiveQueueMessages(
    channel = "orders",
    maxMessages = 1,
    waitTimeoutSeconds = 5
)
for (msg in response.messages) {
    println("Processing: ${String(msg.body)}")
    msg.ack()
    println("Message acknowledged and removed")
}
ack_example.cpp
auto response = client.receiveQueueMessages("orders", 1, 5);
for (const auto& msg : response.messages) {
    std::cout << "Processing: " << msg.body << std::endl;
    msg.ack();
    std::cout << "Message acknowledged and removed" << std::endl;
}
ack_example.rs
let mut receiver = client.new_queue_downstream_receiver().await?;
let response = receiver
    .poll(PollRequest {
        channel: "orders".to_string(),
        max_items: 1,
        wait_timeout_seconds: 5,
        auto_ack: false,
    })
    .await?;
for msg in &response.messages {
    println!("Processing: {}", String::from_utf8_lossy(&msg.message.body));
    msg.ack().await?;
    println!("Message acknowledged and removed");
}
ack_example.rb
receiver = client.create_downstream_receiver
request = KubeMQ::Queues::QueuePollRequest.new(channel: 'orders', max_items: 1, wait_timeout: 5)
response = receiver.poll(request)
response.messages.each do |msg|
  puts "Processing: #{msg.body}"
  msg.ack
  puts 'Message acknowledged and removed'
end
ack_example.exs
{:ok, poll} =
  KubeMQ.Client.poll_queue(client, channel: "orders", max_items: 1, wait_timeout: 5_000)

Enum.each(poll.messages, fn msg -> IO.puts("Processing: #{msg.body}") end)

{:ok, _} = KubeMQ.PollResponse.ack_all(poll)
IO.puts("Messages acknowledged and removed")

Negative Acknowledge (Nack)

Use nack when processing fails and you want the message returned to the queue for another delivery attempt.

nack_example.go
resp, err := client.PollQueue(ctx, &kubemq.PollRequest{
    Channel:            "orders",
    MaxItems:           1,
    WaitTimeoutSeconds: 5,
})
if err != nil {
    log.Fatal(err)
}
for _, m := range resp.Messages {
    fmt.Printf("Attempting: %s\n", string(m.Message.Body))
    // Simulate a processing failure
}
if err := resp.NAckAll(); err != nil {
    log.Fatal(err)
}
fmt.Println("Messages nacked — returned to queue")
nack_example.py
response = client.receive_queue_messages(
    channel="orders",
    max_messages=1,
    wait_timeout_in_seconds=5,
)
for msg in response.messages:
    print(f"Attempting: {msg.body.decode('utf-8')}")
    # Simulate a processing failure
    msg.nack()
    print("Message nacked — returned to queue")
nack_example.ts
const messages = await client.receiveQueueMessages({
  channel: 'orders',
  maxMessages: 1,
  waitTimeoutSeconds: 5,
});
for (const msg of messages) {
  console.log('Attempting:', new TextDecoder().decode(msg.body));
  // Simulate a processing failure
  await msg.nack();
  console.log('Message nacked — returned to queue');
}
NackExample.java
ReceiveQueueMessagesResponse response = client.receiveQueueMessages(
    ReceiveQueueMessagesRequest.builder()
        .channel("orders")
        .maxMessages(1)
        .waitTimeoutSeconds(5)
        .build());

for (QueueMessageReceived msg : response.getMessages()) {
    System.out.println("Attempting: " + new String(msg.getBody()));
    msg.nack();
    System.out.println("Message nacked — returned to queue");
}
NackExample.cs
var response = await client.ReceiveQueueMessagesAsync(new ReceiveQueueMessagesRequest
{
    Channel = "orders",
    MaxMessages = 1,
    WaitTimeoutSeconds = 5,
});
foreach (var msg in response.Messages)
{
    Console.WriteLine($"Attempting: {Encoding.UTF8.GetString(msg.Body.Span)}");
    await msg.NAckAsync();
    Console.WriteLine("Message nacked — returned to queue");
}
NackExample.kt
val response = client.receiveQueueMessages(
    channel = "orders",
    maxMessages = 1,
    waitTimeoutSeconds = 5
)
for (msg in response.messages) {
    println("Attempting: ${String(msg.body)}")
    msg.nack()
    println("Message nacked — returned to queue")
}
nack_example.cpp
auto response = client.receiveQueueMessages("orders", 1, 5);
for (const auto& msg : response.messages) {
    std::cout << "Attempting: " << msg.body << std::endl;
    msg.nack();
    std::cout << "Message nacked — returned to queue" << std::endl;
}
nack_example.rs
let mut receiver = client.new_queue_downstream_receiver().await?;
let response = receiver
    .poll(PollRequest {
        channel: "orders".to_string(),
        max_items: 1,
        wait_timeout_seconds: 5,
        auto_ack: false,
    })
    .await?;
for msg in &response.messages {
    println!("Attempting: {}", String::from_utf8_lossy(&msg.message.body));
    // Simulate a processing failure
    msg.nack().await?;
    println!("Message nacked — returned to queue");
}
nack_example.rb
receiver = client.create_downstream_receiver
request = KubeMQ::Queues::QueuePollRequest.new(channel: 'orders', max_items: 1, wait_timeout: 5)
response = receiver.poll(request)
response.messages.each do |msg|
  puts "Attempting: #{msg.body}"
  # Simulate a processing failure
  msg.nack
  puts 'Message nacked — returned to queue'
end
nack_example.exs
{:ok, poll} =
  KubeMQ.Client.poll_queue(client, channel: "orders", max_items: 1, wait_timeout: 5_000)

Enum.each(poll.messages, fn msg -> IO.puts("Attempting: #{msg.body}") end)

# Simulate a processing failure
{:ok, _} = KubeMQ.PollResponse.nack_all(poll)
IO.puts("Messages nacked — returned to queue")

Requeue to Another Channel

Use requeue to redirect a message to a different queue channel — useful for priority routing, error isolation, or manual review workflows.

requeue_example.go
resp, err := client.PollQueue(ctx, &kubemq.PollRequest{
    Channel:            "orders",
    MaxItems:           1,
    WaitTimeoutSeconds: 5,
})
if err != nil {
    log.Fatal(err)
}
for _, m := range resp.Messages {
    fmt.Printf("Rerouting: %s\n", string(m.Message.Body))
}
if err := resp.ReQueueAll("orders.manual-review"); err != nil {
    log.Fatal(err)
}
fmt.Println("Messages requeued to 'orders.manual-review'")
requeue_example.py
response = client.receive_queue_messages(
    channel="orders",
    max_messages=1,
    wait_timeout_in_seconds=5,
)
for msg in response.messages:
    print(f"Rerouting: {msg.body.decode('utf-8')}")
    msg.requeue("orders.manual-review")
    print("Message requeued to 'orders.manual-review'")
requeue_example.ts
const messages = await client.receiveQueueMessages({
  channel: 'orders',
  maxMessages: 1,
  waitTimeoutSeconds: 5,
});
for (const msg of messages) {
  console.log('Rerouting:', new TextDecoder().decode(msg.body));
  await msg.requeue('orders.manual-review');
  console.log("Message requeued to 'orders.manual-review'");
}
RequeueExample.java
ReceiveQueueMessagesResponse response = client.receiveQueueMessages(
    ReceiveQueueMessagesRequest.builder()
        .channel("orders")
        .maxMessages(1)
        .waitTimeoutSeconds(5)
        .build());

for (QueueMessageReceived msg : response.getMessages()) {
    System.out.println("Rerouting: " + new String(msg.getBody()));
    msg.requeue("orders.manual-review");
    System.out.println("Message requeued to 'orders.manual-review'");
}
RequeueExample.cs
var response = await client.ReceiveQueueMessagesAsync(new ReceiveQueueMessagesRequest
{
    Channel = "orders",
    MaxMessages = 1,
    WaitTimeoutSeconds = 5,
});
foreach (var msg in response.Messages)
{
    Console.WriteLine($"Rerouting: {Encoding.UTF8.GetString(msg.Body.Span)}");
    await msg.ReQueueAsync("orders.manual-review");
    Console.WriteLine("Message requeued to 'orders.manual-review'");
}
RequeueExample.kt
val response = client.receiveQueueMessages(
    channel = "orders",
    maxMessages = 1,
    waitTimeoutSeconds = 5
)
for (msg in response.messages) {
    println("Rerouting: ${String(msg.body)}")
    msg.requeue("orders.manual-review")
    println("Message requeued to 'orders.manual-review'")
}
requeue_example.cpp
auto response = client.receiveQueueMessages("orders", 1, 5);
for (const auto& msg : response.messages) {
    std::cout << "Rerouting: " << msg.body << std::endl;
    msg.requeue("orders.manual-review");
    std::cout << "Message requeued to 'orders.manual-review'" << std::endl;
}
requeue_example.rs
let mut receiver = client.new_queue_downstream_receiver().await?;
let response = receiver
    .poll(PollRequest {
        channel: "orders".to_string(),
        max_items: 1,
        wait_timeout_seconds: 5,
        auto_ack: false,
    })
    .await?;
for msg in &response.messages {
    println!("Rerouting: {}", String::from_utf8_lossy(&msg.message.body));
    msg.re_queue("orders.manual-review").await?;
    println!("Message requeued to 'orders.manual-review'");
}
requeue_example.rb
receiver = client.create_downstream_receiver
request = KubeMQ::Queues::QueuePollRequest.new(channel: 'orders', max_items: 1, wait_timeout: 5)
response = receiver.poll(request)
response.messages.each do |msg|
  puts "Rerouting: #{msg.body}"
end
# Requeue is settled at the poll-response level, not per message
response.requeue_all(channel: 'orders.manual-review')
puts "Messages requeued to 'orders.manual-review'"
requeue_example.exs
{:ok, poll} =
  KubeMQ.Client.poll_queue(client, channel: "orders", max_items: 1, wait_timeout: 5_000)

Enum.each(poll.messages, fn msg -> IO.puts("Rerouting: #{msg.body}") end)

# Elixir settles a poll transaction as a whole — requeue_all moves them to the target channel
{:ok, _} = KubeMQ.PollResponse.requeue_all(poll, "orders.manual-review")
IO.puts("Messages requeued to 'orders.manual-review'")

When to Use Each Option

OptionEffectUse When
AckRemove from queue permanentlyProcessing succeeded
NackReturn to same queue for redeliveryTransient failure (network, timeout)
RequeueMove to a different queue channelNeeds manual review, priority routing, or error isolation

If a message is nacked repeatedly and exceeds the maxReceiveCount, it is automatically routed to the dead letter queue (if configured). See Dead Letter Queue for details.

Next Steps

Was this page helpful?

On this page