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
- KubeMQ server running on
localhost:50000 - SDK installed (Getting Started)
Steps
Send a Delayed Message
Set the delay in seconds when creating the message. The message becomes available after the delay expires.
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)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}")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}`);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());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}");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}")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;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
);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}"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.
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()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()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();
}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();
}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();
}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()
}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();
}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)
);
}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}"
endIO.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:
| Delay | Expiration | Message Available | Message Expires |
|---|---|---|---|
| 30s | 0 | T+30s | Never |
| 0 | 60s | Immediately | T+60s |
| 30s | 60s | T+30s | T+90s |
The maximum delay is controlled by the server setting MaxDelaySeconds (default: 43,200 seconds / 12 hours).
Next Steps
Was this page helpful?