KubeMQ
LearnQueuesTutorials

Delayed Messages

Schedule messages for future delivery using KubeMQ delayed message queues.

How Delayed Messages Work

Messages sent with delaySeconds > 0 are held in an internal delay channel (_QUEUE_DELAY_). A background processor checks every 500ms for expired delays and moves messages to the target queue.

A delayed message waits in the internal _QUEUE_DELAY_ channel until its delay expires, then moves to the target queue for normal consumption.

Prerequisites

Steps

Send a Delayed Message

Set the delay in seconds when creating the message. The message becomes available after the delay expires.

delayed_sender.go
msg := kubemq.NewQueueMessage().
    SetChannel("scheduled-tasks").
    SetBody([]byte(`{"task":"send-reminder","orderId":"ORD-5001"}`)).
    SetDelaySeconds(30)

result, err := client.SendQueueMessage(ctx, msg)
if err != nil {
    log.Fatal(err)
}
fmt.Printf("Sent delayed message: id=%s, delayedTo=%d\n",
    result.MessageID, result.DelayedTo)
delayed_sender.py
result = client.send_queue_message(
    QueueMessage(
        channel="scheduled-tasks",
        body=b'{"task":"send-reminder","orderId":"ORD-5001"}',
        delay_in_seconds=30,
    )
)
print(f"Sent delayed message: id={result.id}, delayedTo={result.delayed_to}")
delayed_sender.ts
const result = await client.sendQueueMessage(
  createQueueMessage({
    channel: 'scheduled-tasks',
    body: JSON.stringify({ task: 'send-reminder', orderId: 'ORD-5001' }),
    policy: { delaySeconds: 30 },
  }),
);
console.log(`Sent delayed message: id=${result.messageId}, delayedTo=${result.delayedTo}`);
DelayedSender.java
QueueMessage msg = QueueMessage.builder()
    .channel("scheduled-tasks")
    .body("{\"task\":\"send-reminder\",\"orderId\":\"ORD-5001\"}".getBytes())
    .delaySeconds(30)
    .build();

SendQueueMessageResult result = client.sendQueueMessage(msg);
System.out.printf("Sent delayed message: id=%s, delayedTo=%d%n",
    result.getMessageId(), result.getDelayedTo());
DelayedSender.cs
var result = await client.SendQueueMessageAsync(new QueueMessage
{
    Channel = "scheduled-tasks",
    Body = Encoding.UTF8.GetBytes("{\"task\":\"send-reminder\",\"orderId\":\"ORD-5001\"}"),
    DelaySeconds = 30
});
Console.WriteLine($"Sent delayed message: id={result.MessageId}, delayedTo={result.DelayedTo}");
DelayedSender.kt
val result = client.sendQueueMessage(QueueMessage(
    channel = "scheduled-tasks",
    body = """{"task":"send-reminder","orderId":"ORD-5001"}""".toByteArray(),
    delaySeconds = 30
))
println("Sent delayed message: id=${result.messageId}, delayedTo=${result.delayedTo}")
delayed_sender.cpp
kubemq::QueueMessage msg;
msg.channel = "scheduled-tasks";
msg.body = R"({"task":"send-reminder","orderId":"ORD-5001"})";
msg.delaySeconds = 30;

auto result = client.sendQueueMessage(msg);
std::cout << "Sent delayed message: id=" << result.messageId << std::endl;
delayed_sender.rs
let msg = QueueMessageBuilder::new()
    .channel("scheduled-tasks")
    .body(br#"{"task":"send-reminder","orderId":"ORD-5001"}"#.to_vec())
    .delay_seconds(30)
    .build();

let result = client.send_queue_message(msg).await?;
println!(
    "Sent delayed message: id={}, delayed_to={}",
    result.message_id, result.delayed_to
);
delayed_sender.rb
policy = KubeMQ::Queues::QueueMessagePolicy.new(delay_seconds: 30)
msg = KubeMQ::Queues::QueueMessage.new(
  channel: 'scheduled-tasks',
  body: '{"task":"send-reminder","orderId":"ORD-5001"}',
  policy: policy
)

result = client.send_queue_message(msg)
puts "Sent delayed message: id=#{result.id}, delayed_to=#{result.delayed_to}"
delayed_sender.exs
msg = KubeMQ.QueueMessage.new(
  channel: "scheduled-tasks",
  body: ~s({"task":"send-reminder","orderId":"ORD-5001"}),
  policy: KubeMQ.QueuePolicy.new(delay_seconds: 30)
)

{:ok, result} = KubeMQ.Client.send_queue_message(client, msg)
IO.puts("Sent delayed message, delayed_to: #{result.delayed_to}")

Receive After Delay

Poll the queue after the delay expires. The message is not visible until the delay elapses.

delayed_receiver.go
fmt.Println("Waiting for delayed message...")
time.Sleep(35 * time.Second)

resp, err := client.PollQueue(ctx, &kubemq.PollRequest{
    Channel:            "scheduled-tasks",
    MaxItems:           1,
    WaitTimeoutSeconds: 5,
})
if err != nil {
    log.Fatal(err)
}
for _, m := range resp.Messages {
    fmt.Printf("Received delayed message: %s\n", string(m.Message.Body))
}
resp.AckAll()
delayed_receiver.py
import time

print("Waiting for delayed message...")
time.sleep(35)

response = client.receive_queue_messages(
    channel="scheduled-tasks",
    max_messages=1,
    wait_timeout_in_seconds=5,
)
for msg in response.messages:
    print(f"Received delayed message: {msg.body.decode('utf-8')}")
    msg.ack()
delayed_receiver.ts
console.log('Waiting for delayed message...');
await new Promise((r) => setTimeout(r, 35000));

const messages = await client.receiveQueueMessages({
  channel: 'scheduled-tasks',
  maxMessages: 1,
  waitTimeoutSeconds: 5,
});
for (const msg of messages) {
  console.log('Received delayed message:', new TextDecoder().decode(msg.body));
  await msg.ack();
}
DelayedReceiver.java
System.out.println("Waiting for delayed message...");
Thread.sleep(35000);

ReceiveQueueMessagesResponse response = client.receiveQueueMessages(
    ReceiveQueueMessagesRequest.builder()
        .channel("scheduled-tasks")
        .maxMessages(1)
        .waitTimeoutSeconds(5)
        .build());

for (QueueMessageReceived msg : response.getMessages()) {
    System.out.println("Received delayed message: " + new String(msg.getBody()));
    msg.ack();
}
DelayedReceiver.cs
Console.WriteLine("Waiting for delayed message...");
await Task.Delay(35000);

var response = await client.ReceiveQueueMessagesAsync(new ReceiveQueueMessagesRequest
{
    Channel = "scheduled-tasks",
    MaxMessages = 1,
    WaitTimeoutSeconds = 5,
});
foreach (var msg in response.Messages)
{
    Console.WriteLine($"Received delayed message: {Encoding.UTF8.GetString(msg.Body.Span)}");
    await msg.AckAsync();
}
DelayedReceiver.kt
println("Waiting for delayed message...")
Thread.sleep(35000)

val response = client.receiveQueueMessages(
    channel = "scheduled-tasks",
    maxMessages = 1,
    waitTimeoutSeconds = 5
)
for (msg in response.messages) {
    println("Received delayed message: ${String(msg.body)}")
    msg.ack()
}
delayed_receiver.cpp
std::cout << "Waiting for delayed message..." << std::endl;
std::this_thread::sleep_for(std::chrono::seconds(35));

auto response = client.receiveQueueMessages("scheduled-tasks", 1, 5);
for (const auto& msg : response.messages) {
    std::cout << "Received delayed message: " << msg.body << std::endl;
    msg.ack();
}
delayed_receiver.rs
println!("Waiting for delayed message...");
tokio::time::sleep(std::time::Duration::from_secs(35)).await;

// receive_queue_messages(channel, max_messages, wait_time_seconds, is_peek)
let messages = client
    .receive_queue_messages("scheduled-tasks", 1, 5, false)
    .await?;
for msg in &messages {
    println!(
        "Received delayed message: {}",
        String::from_utf8_lossy(&msg.body)
    );
}
delayed_receiver.rb
puts 'Waiting for delayed message...'
sleep 35

received = client.receive_queue_messages(
  channel: 'scheduled-tasks',
  max_messages: 1,
  wait_timeout_seconds: 5
)
received.each do |msg|
  puts "Received delayed message: #{msg.body}"
end
delayed_receiver.exs
IO.puts("Waiting for delayed message...")
Process.sleep(35_000)

{:ok, result} =
  KubeMQ.Client.receive_queue_messages(client, "scheduled-tasks",
    max_messages: 1,
    wait_timeout: 5_000
  )

Enum.each(result.messages, fn msg ->
  IO.puts("Received delayed message: #{msg.body}")
end)

Delay + Expiration Interaction

When both delaySeconds and expirationSeconds are set, the expiration clock starts after the delay expires:

DelayExpirationMessage AvailableMessage Expires
30s0T+30sNever
060sImmediatelyT+60s
30s60sT+30sT+90s

The maximum delay is controlled by the server setting MaxDelaySeconds (default: 43,200 seconds / 12 hours).

Next Steps

Was this page helpful?

On this page