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
- KubeMQ server running on
localhost:50000 - SDK installed (Getting Started)
Steps
Send a Message with DLQ Policy
Set maxReceiveCount and maxReceiveQueue on the message to enable dead letter routing.
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)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}")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}`);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());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}");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}")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;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);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'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.
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")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")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');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");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");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")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;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");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'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.
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()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()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();
}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();
}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();
}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()
}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();
}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_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}"
endcase 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")
endDLQ Message Attributes
When a message is routed to the dead letter queue:
| Attribute | Value |
|---|---|
reRouted | true |
reRoutedFromQueue | Original source channel name |
receiveCount | Reset to 0 in the DLQ |
| Policy fields | Reset 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?