KubeMQ
LearnQueuesTutorials

Dead Letter Queue

Route failed messages to a dead letter queue after exceeding the maximum receive count.

How DLQ Works

When a message is nacked or its visibility timeout expires repeatedly, the receiveCount increments. Once it exceeds maxReceiveCount, KubeMQ automatically routes the message to the dead letter queue channel.

Messages that exhaust maxReceiveCount are routed from the source queue to the dead letter queue for inspection and recovery.

Prerequisites

Steps

Send a Message with DLQ Policy

Set maxReceiveCount and maxReceiveQueue on the message to enable dead letter routing.

dlq_sender.go
msg := kubemq.NewQueueMessage().
    SetChannel("orders").
    SetBody([]byte(`{"orderId":"ORD-9001","total":250.00}`)).
    SetMaxReceiveCount(3).
    SetMaxReceiveQueue("orders.dlq")

result, err := client.SendQueueMessage(ctx, msg)
if err != nil {
    log.Fatal(err)
}
fmt.Printf("Sent with DLQ policy: id=%s\n", result.MessageID)
dlq_sender.py
result = client.send_queue_message(
    QueueMessage(
        channel="orders",
        body=b'{"orderId":"ORD-9001","total":250.00}',
        max_receive_count=3,
        max_receive_queue="orders.dlq",
    )
)
print(f"Sent with DLQ policy: id={result.id}")
dlq_sender.ts
const result = await client.sendQueueMessage(
  createQueueMessage({
    channel: 'orders',
    body: JSON.stringify({ orderId: 'ORD-9001', total: 250.0 }),
    policy: {
      maxReceiveCount: 3,
      maxReceiveQueue: 'orders.dlq',
    },
  }),
);
console.log(`Sent with DLQ policy: id=${result.messageId}`);
DlqSender.java
QueueMessage msg = QueueMessage.builder()
    .channel("orders")
    .body("{\"orderId\":\"ORD-9001\",\"total\":250.00}".getBytes())
    .maxReceiveCount(3)
    .maxReceiveQueue("orders.dlq")
    .build();

SendQueueMessageResult result = client.sendQueueMessage(msg);
System.out.printf("Sent with DLQ policy: id=%s%n", result.getMessageId());
DlqSender.cs
var result = await client.SendQueueMessageAsync(new QueueMessage
{
    Channel = "orders",
    Body = Encoding.UTF8.GetBytes("{\"orderId\":\"ORD-9001\",\"total\":250.00}"),
    MaxReceiveCount = 3,
    MaxReceiveQueue = "orders.dlq"
});
Console.WriteLine($"Sent with DLQ policy: id={result.MessageId}");
DlqSender.kt
val result = client.sendQueueMessage(QueueMessage(
    channel = "orders",
    body = """{"orderId":"ORD-9001","total":250.00}""".toByteArray(),
    maxReceiveCount = 3,
    maxReceiveQueue = "orders.dlq"
))
println("Sent with DLQ policy: id=${result.messageId}")
dlq_sender.cpp
kubemq::QueueMessage msg;
msg.channel = "orders";
msg.body = R"({"orderId":"ORD-9001","total":250.00})";
msg.maxReceiveCount = 3;
msg.maxReceiveQueue = "orders.dlq";

auto result = client.sendQueueMessage(msg);
std::cout << "Sent with DLQ policy: id=" << result.messageId << std::endl;
dlq_sender.rs
let msg = QueueMessageBuilder::new()
    .channel("orders")
    .body(br#"{"orderId":"ORD-9001","total":250.00}"#.to_vec())
    .max_receive_count(3)
    .max_receive_queue("orders.dlq")
    .build();

let result = client.send_queue_message(msg).await?;
println!("Sent with DLQ policy: id={}", result.message_id);
dlq_sender.rb
policy = KubeMQ::Queues::QueueMessagePolicy.new(
  max_receive_count: 3,
  max_receive_queue: 'orders.dlq'
)
msg = KubeMQ::Queues::QueueMessage.new(
  channel: 'orders',
  body: '{"orderId":"ORD-9001","total":250.00}',
  policy: policy
)

client.send_queue_message(msg)
puts 'Sent with DLQ policy: orders.dlq'
dlq_sender.exs
msg = KubeMQ.QueueMessage.new(
  channel: "orders",
  body: ~s({"orderId":"ORD-9001","total":250.00}),
  policy: KubeMQ.QueuePolicy.new(
    max_receive_count: 3,
    max_receive_queue: "orders.dlq"
  )
)

{:ok, result} = KubeMQ.Client.send_queue_message(client, msg)
IO.puts("Sent with DLQ policy: id=#{result.message_id}")

Simulate Failures

Receive and nack the message 3 times to trigger DLQ routing.

dlq_nacker.go
for attempt := 1; attempt <= 3; attempt++ {
    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("Attempt %d: receiveCount=%d\n", attempt, m.Message.Attributes.ReceiveCount)
    }
    resp.NAckAll()
    time.Sleep(time.Second)
}
fmt.Println("Message should now be in DLQ")
dlq_nacker.py
import time

for attempt in range(1, 4):
    response = client.receive_queue_messages(
        channel="orders",
        max_messages=1,
        wait_timeout_in_seconds=5,
    )
    for msg in response.messages:
        print(f"Attempt {attempt}: receiveCount={msg.receive_count}")
        msg.nack()
    time.sleep(1)

print("Message should now be in DLQ")
dlq_nacker.ts
for (let attempt = 1; attempt <= 3; attempt++) {
  const messages = await client.receiveQueueMessages({
    channel: 'orders',
    maxMessages: 1,
    waitTimeoutSeconds: 5,
  });
  for (const msg of messages) {
    console.log(`Attempt ${attempt}: receiveCount=${msg.receiveCount}`);
    await msg.nack();
  }
  await new Promise((r) => setTimeout(r, 1000));
}
console.log('Message should now be in DLQ');
DlqNacker.java
for (int attempt = 1; attempt <= 3; attempt++) {
    ReceiveQueueMessagesResponse response = client.receiveQueueMessages(
        ReceiveQueueMessagesRequest.builder()
            .channel("orders")
            .maxMessages(1)
            .waitTimeoutSeconds(5)
            .build());

    for (QueueMessageReceived msg : response.getMessages()) {
        System.out.printf("Attempt %d: receiveCount=%d%n", attempt, msg.getReceiveCount());
        msg.nack();
    }
    Thread.sleep(1000);
}
System.out.println("Message should now be in DLQ");
DlqNacker.cs
for (int attempt = 1; attempt <= 3; attempt++)
{
    var response = await client.ReceiveQueueMessagesAsync(new ReceiveQueueMessagesRequest
    {
        Channel = "orders",
        MaxMessages = 1,
        WaitTimeoutSeconds = 5,
    });
    foreach (var msg in response.Messages)
    {
        Console.WriteLine($"Attempt {attempt}: receiveCount={msg.ReceiveCount}");
        await msg.NAckAsync();
    }
    await Task.Delay(1000);
}
Console.WriteLine("Message should now be in DLQ");
DlqNacker.kt
for (attempt in 1..3) {
    val response = client.receiveQueueMessages(
        channel = "orders",
        maxMessages = 1,
        waitTimeoutSeconds = 5
    )
    for (msg in response.messages) {
        println("Attempt $attempt: receiveCount=${msg.receiveCount}")
        msg.nack()
    }
    Thread.sleep(1000)
}
println("Message should now be in DLQ")
dlq_nacker.cpp
for (int attempt = 1; attempt <= 3; attempt++) {
    auto response = client.receiveQueueMessages("orders", 1, 5);
    for (const auto& msg : response.messages) {
        std::cout << "Attempt " << attempt
                  << ": receiveCount=" << msg.receiveCount << std::endl;
        msg.nack();
    }
    std::this_thread::sleep_for(std::chrono::seconds(1));
}
std::cout << "Message should now be in DLQ" << std::endl;
dlq_nacker.rs
let mut receiver = client.new_queue_downstream_receiver().await?;
for attempt in 1..=3 {
    let poll = PollRequest {
        channel: "orders".to_string(),
        max_items: 1,
        wait_timeout_seconds: 5,
        auto_ack: false,
    };
    let response = receiver.poll(poll).await?;
    println!("Attempt {}: {} messages", attempt, response.messages.len());
    if !response.messages.is_empty() {
        response.nack_all().await?;
    }
    tokio::time::sleep(std::time::Duration::from_secs(1)).await;
}
println!("Message should now be in DLQ");
dlq_nacker.rb
receiver = client.create_downstream_receiver

(1..3).each do |attempt|
  request = KubeMQ::Queues::QueuePollRequest.new(
    channel: 'orders',
    max_items: 1,
    wait_timeout: 5
  )
  response = receiver.poll(request)
  puts "Attempt #{attempt}: #{response.messages.size} messages"
  response.nack_all if response.messages.any?
  sleep 1
end
puts 'Message should now be in DLQ'
dlq_nacker.exs
for attempt <- 1..3 do
  case KubeMQ.Client.poll_queue(client,
         channel: "orders",
         max_items: 1,
         wait_timeout: 5_000
       ) do
    {:ok, poll} when length(poll.messages) > 0 ->
      IO.puts("Attempt #{attempt}: #{length(poll.messages)} messages")
      :ok = KubeMQ.PollResponse.nack_all(poll)

    _ ->
      IO.puts("Attempt #{attempt}: no message")
  end

  Process.sleep(1_000)
end

IO.puts("Message should now be in DLQ")

Read from the Dead Letter Queue

The failed message now sits in the DLQ channel with routing metadata attached.

dlq_reader.go
resp, err := client.PollQueue(ctx, &kubemq.PollRequest{
    Channel:            "orders.dlq",
    MaxItems:           10,
    WaitTimeoutSeconds: 5,
})
if err != nil {
    log.Fatal(err)
}
for _, m := range resp.Messages {
    fmt.Printf("DLQ message: %s\n", string(m.Message.Body))
    fmt.Printf("  Rerouted from: %s\n", m.Message.Attributes.ReRoutedFromQueue)
    fmt.Printf("  Receive count: %d\n", m.Message.Attributes.ReceiveCount)
}
resp.AckAll()
dlq_reader.py
response = client.receive_queue_messages(
    channel="orders.dlq",
    max_messages=10,
    wait_timeout_in_seconds=5,
)
for msg in response.messages:
    print(f"DLQ message: {msg.body.decode('utf-8')}")
    print(f"  Rerouted from: {msg.rerouted_from_queue}")
    print(f"  Receive count: {msg.receive_count}")
    msg.ack()
dlq_reader.ts
const dlqMessages = await client.receiveQueueMessages({
  channel: 'orders.dlq',
  maxMessages: 10,
  waitTimeoutSeconds: 5,
});
for (const msg of dlqMessages) {
  console.log('DLQ message:', new TextDecoder().decode(msg.body));
  console.log('  Rerouted from:', msg.reRoutedFromQueue);
  console.log('  Receive count:', msg.receiveCount);
  await msg.ack();
}
DlqReader.java
ReceiveQueueMessagesResponse dlqResponse = client.receiveQueueMessages(
    ReceiveQueueMessagesRequest.builder()
        .channel("orders.dlq")
        .maxMessages(10)
        .waitTimeoutSeconds(5)
        .build());

for (QueueMessageReceived msg : dlqResponse.getMessages()) {
    System.out.println("DLQ message: " + new String(msg.getBody()));
    System.out.println("  Rerouted from: " + msg.getReRoutedFromQueue());
    System.out.println("  Receive count: " + msg.getReceiveCount());
    msg.ack();
}
DlqReader.cs
var dlqResponse = await client.ReceiveQueueMessagesAsync(new ReceiveQueueMessagesRequest
{
    Channel = "orders.dlq",
    MaxMessages = 10,
    WaitTimeoutSeconds = 5,
});
foreach (var msg in dlqResponse.Messages)
{
    Console.WriteLine($"DLQ message: {Encoding.UTF8.GetString(msg.Body.Span)}");
    Console.WriteLine($"  Rerouted from: {msg.ReRoutedFromQueue}");
    Console.WriteLine($"  Receive count: {msg.ReceiveCount}");
    await msg.AckAsync();
}
DlqReader.kt
val dlqResponse = client.receiveQueueMessages(
    channel = "orders.dlq",
    maxMessages = 10,
    waitTimeoutSeconds = 5
)
for (msg in dlqResponse.messages) {
    println("DLQ message: ${String(msg.body)}")
    println("  Rerouted from: ${msg.reRoutedFromQueue}")
    println("  Receive count: ${msg.receiveCount}")
    msg.ack()
}
dlq_reader.cpp
auto dlqResponse = client.receiveQueueMessages("orders.dlq", 10, 5);
for (const auto& msg : dlqResponse.messages) {
    std::cout << "DLQ message: " << msg.body << std::endl;
    std::cout << "  Rerouted from: " << msg.reRoutedFromQueue << std::endl;
    std::cout << "  Receive count: " << msg.receiveCount << std::endl;
    msg.ack();
}
dlq_reader.rs
let dlq_msgs = client
    .receive_queue_messages("orders.dlq", 10, 5, false)
    .await?;
for msg in &dlq_msgs {
    println!("DLQ message: {}", String::from_utf8_lossy(&msg.body));
    if let Some(attr) = &msg.attributes {
        println!("  Rerouted from: {}", attr.re_routed_from_queue);
        println!("  Receive count: {}", attr.receive_count);
    }
}
dlq_reader.rb
dlq_messages = client.receive_queue_messages(
  channel: 'orders.dlq',
  max_messages: 10,
  wait_timeout_seconds: 5
)
dlq_messages.each do |msg|
  puts "DLQ message: #{msg.body}"
  puts "  Rerouted from: #{msg.attributes.re_routed_from_queue}"
  puts "  Receive count: #{msg.attributes.receive_count}"
end
dlq_reader.exs
case KubeMQ.Client.receive_queue_messages(client, "orders.dlq",
       max_messages: 10,
       wait_timeout: 5_000
     ) do
  {:ok, result} ->
    Enum.each(result.messages, fn msg ->
      IO.puts("DLQ message: #{msg.body}")

      if msg.attributes do
        IO.puts("  Rerouted from: #{msg.attributes.re_routed_from_queue}")
        IO.puts("  Receive count: #{msg.attributes.receive_count}")
      end
    end)

  {:error, _} ->
    IO.puts("No DLQ messages yet")
end

DLQ Message Attributes

When a message is routed to the dead letter queue:

AttributeValue
reRoutedtrue
reRoutedFromQueueOriginal source channel name
receiveCountReset to 0 in the DLQ
Policy fieldsReset to server defaults

If maxReceiveQueue is not set on the message, messages that exceed maxReceiveCount are silently discarded. Always set a DLQ channel for critical workloads.

Next Steps

Was this page helpful?

On this page