Scheduled Task System
Implement cron-like scheduled tasks using KubeMQ delayed messages.
Architecture
A scheduler submits tasks with a delay — the task becomes available only after the delay expires. Task workers process tasks when they become due. For recurring tasks, the worker re-enqueues the task with the next delay.
A scheduler sends each task with a delay; the broker holds it until the delay expires, then a worker receives and processes it — re-enqueuing recurring tasks for their next run.
Schedule a Task
Send a message with delaySeconds set to the desired future execution time.
type ScheduledTask struct {
TaskID string `json:"taskId"`
Action string `json:"action"`
Recurring bool `json:"recurring"`
IntervalSec int `json:"intervalSec"`
}
func scheduleTask(ctx context.Context, client *kubemq.Client, task ScheduledTask, delaySec int) error {
body, _ := json.Marshal(task)
msg := kubemq.NewQueueMessage().
SetChannel("scheduled-tasks").
SetBody(body).
SetMetadata(task.Action).
SetTags(map[string]string{
"taskId": task.TaskID,
"action": task.Action,
}).
SetDelaySeconds(delaySec)
result, err := client.SendQueueMessage(ctx, msg)
if err != nil {
return err
}
log.Printf("Task %s scheduled in %ds: id=%s", task.TaskID, delaySec, result.MessageID)
return nil
}import json
def schedule_task(client, task, delay_sec):
result = client.send_queue_message(
QueueMessage(
channel="scheduled-tasks",
body=json.dumps(task).encode(),
metadata=task["action"],
tags={"taskId": task["taskId"], "action": task["action"]},
delay_in_seconds=delay_sec,
)
)
print(f"Task {task['taskId']} scheduled in {delay_sec}s: id={result.id}")
schedule_task(client, {
"taskId": "reminder-001",
"action": "send-reminder-email",
"recurring": True,
"intervalSec": 3600,
}, delay_sec=3600)async function scheduleTask(client: KubeMQClient, task: ScheduledTask, delaySec: number) {
const result = await client.sendQueueMessage(
createQueueMessage({
channel: 'scheduled-tasks',
body: JSON.stringify(task),
metadata: task.action,
tags: { taskId: task.taskId, action: task.action },
policy: { delaySeconds: delaySec },
}),
);
console.log(`Task ${task.taskId} scheduled in ${delaySec}s: id=${result.messageId}`);
}public void scheduleTask(QueuesClient client, ScheduledTask task, int delaySec) throws Exception {
QueueMessage msg = QueueMessage.builder()
.channel("scheduled-tasks")
.body(objectMapper.writeValueAsBytes(task))
.metadata(task.getAction())
.tags(Map.of("taskId", task.getTaskId(), "action", task.getAction()))
.delaySeconds(delaySec)
.build();
SendQueueMessageResult result = client.sendQueueMessage(msg);
System.out.printf("Task %s scheduled in %ds: id=%s%n",
task.getTaskId(), delaySec, result.getMessageId());
}async Task ScheduleTask(QueuesClient client, ScheduledTask task, int delaySec)
{
var result = await client.SendQueueMessageAsync(new QueueMessage
{
Channel = "scheduled-tasks",
Body = JsonSerializer.SerializeToUtf8Bytes(task),
Metadata = task.Action,
Tags = new Dictionary<string, string>
{
["taskId"] = task.TaskId,
["action"] = task.Action
},
DelaySeconds = delaySec
});
Console.WriteLine($"Task {task.TaskId} scheduled in {delaySec}s: id={result.MessageId}");
}suspend fun scheduleTask(client: QueuesClient, task: ScheduledTask, delaySec: Int) {
val result = client.sendQueueMessage(QueueMessage(
channel = "scheduled-tasks",
body = Json.encodeToString(task).toByteArray(),
metadata = task.action,
tags = mapOf("taskId" to task.taskId, "action" to task.action),
delaySeconds = delaySec
))
println("Task ${task.taskId} scheduled in ${delaySec}s: id=${result.messageId}")
}void scheduleTask(kubemq::QueuesClient& client, const ScheduledTask& task, int delaySec) {
kubemq::QueueMessage msg;
msg.channel = "scheduled-tasks";
msg.body = task.toJson();
msg.metadata = task.action;
msg.tags = {{"taskId", task.taskId}, {"action", task.action}};
msg.delaySeconds = delaySec;
auto result = client.sendQueueMessage(msg);
std::cout << "Task " << task.taskId << " scheduled in "
<< delaySec << "s: id=" << result.messageId << std::endl;
}use kubemq::prelude::*;
use kubemq::QueueMessageBuilder;
use std::collections::HashMap;
async fn schedule_task(
client: &KubemqClient,
task: &ScheduledTask,
delay_sec: i32,
) -> kubemq::Result<()> {
let mut tags = HashMap::new();
tags.insert("taskId".to_string(), task.task_id.clone());
tags.insert("action".to_string(), task.action.clone());
let msg = QueueMessageBuilder::new()
.channel("scheduled-tasks")
.body(serde_json::to_vec(task).unwrap())
.metadata(task.action.clone())
.tags(tags)
.delay_seconds(delay_sec)
.build();
let result = client.send_queue_message(msg).await?;
println!(
"Task {} scheduled in {}s: id={}",
task.task_id, delay_sec, result.message_id
);
Ok(())
}require 'kubemq'
require 'json'
def schedule_task(client, task, delay_sec)
policy = KubeMQ::Queues::QueueMessagePolicy.new(delay_seconds: delay_sec)
msg = KubeMQ::Queues::QueueMessage.new(
channel: 'scheduled-tasks',
metadata: task['action'],
body: task.to_json,
tags: { 'taskId' => task['taskId'], 'action' => task['action'] },
policy: policy
)
result = client.send_queue_message(msg)
puts "Task #{task['taskId']} scheduled in #{delay_sec}s: id=#{result.id}"
end
schedule_task(client, {
'taskId' => 'reminder-001',
'action' => 'send-reminder-email',
'recurring' => true,
'intervalSec' => 3600
}, 3600)defp schedule_task(client, task, delay_sec) do
msg =
KubeMQ.QueueMessage.new(
channel: "scheduled-tasks",
metadata: task.action,
body: Jason.encode!(task),
tags: %{"taskId" => task.task_id, "action" => task.action},
policy: KubeMQ.QueuePolicy.new(delay_seconds: delay_sec)
)
{:ok, result} = KubeMQ.Client.send_queue_message(client, msg)
IO.puts("Task #{task.task_id} scheduled in #{delay_sec}s: id=#{result.message_id}")
endTask Worker with Recurring Support
When a task completes, check if it's recurring and re-schedule it.
for {
resp, err := client.PollQueue(ctx, &kubemq.PollRequest{
Channel: "scheduled-tasks",
MaxItems: 5,
WaitTimeoutSeconds: 10,
})
if err != nil {
log.Printf("Poll error: %v", err)
time.Sleep(5 * time.Second)
continue
}
for _, m := range resp.Messages {
var task ScheduledTask
json.Unmarshal(m.Message.Body, &task)
log.Printf("Executing task %s: %s", task.TaskID, task.Action)
executeTask(task)
if task.Recurring {
scheduleTask(ctx, client, task, task.IntervalSec)
log.Printf("Task %s rescheduled in %ds", task.TaskID, task.IntervalSec)
}
}
resp.AckAll()
}import json
import time
while True:
response = client.receive_queue_messages(
channel="scheduled-tasks",
max_messages=5,
wait_timeout_in_seconds=10,
)
for msg in response.messages:
task = json.loads(msg.body.decode("utf-8"))
print(f"Executing task {task['taskId']}: {task['action']}")
execute_task(task)
if task.get("recurring"):
schedule_task(client, task, task["intervalSec"])
print(f"Task {task['taskId']} rescheduled in {task['intervalSec']}s")
msg.ack()
time.sleep(1)while (true) {
const messages = await client.receiveQueueMessages({
channel: 'scheduled-tasks',
maxMessages: 5,
waitTimeoutSeconds: 10,
});
for (const msg of messages) {
const task = JSON.parse(new TextDecoder().decode(msg.body));
console.log(`Executing task ${task.taskId}: ${task.action}`);
await executeTask(task);
if (task.recurring) {
await scheduleTask(client, task, task.intervalSec);
console.log(`Task ${task.taskId} rescheduled in ${task.intervalSec}s`);
}
await msg.ack();
}
await new Promise((r) => setTimeout(r, 1000));
}while (true) {
ReceiveQueueMessagesResponse response = client.receiveQueueMessages(
ReceiveQueueMessagesRequest.builder()
.channel("scheduled-tasks").maxMessages(5)
.waitTimeoutSeconds(10).build());
for (QueueMessageReceived msg : response.getMessages()) {
ScheduledTask task = objectMapper.readValue(msg.getBody(), ScheduledTask.class);
System.out.printf("Executing task %s: %s%n", task.getTaskId(), task.getAction());
executeTask(task);
if (task.isRecurring()) {
scheduleTask(client, task, task.getIntervalSec());
}
msg.ack();
}
Thread.sleep(1000);
}while (true)
{
var response = await client.ReceiveQueueMessagesAsync(new ReceiveQueueMessagesRequest
{
Channel = "scheduled-tasks", MaxMessages = 5, WaitTimeoutSeconds = 10,
});
foreach (var msg in response.Messages)
{
var task = JsonSerializer.Deserialize<ScheduledTask>(msg.Body.Span);
Console.WriteLine($"Executing task {task.TaskId}: {task.Action}");
ExecuteTask(task);
if (task.Recurring)
await ScheduleTask(client, task, task.IntervalSec);
await msg.AckAsync();
}
await Task.Delay(1000);
}while (true) {
val response = client.receiveQueueMessages(
channel = "scheduled-tasks", maxMessages = 5, waitTimeoutSeconds = 10)
for (msg in response.messages) {
val task = Json.decodeFromString<ScheduledTask>(String(msg.body))
println("Executing task ${task.taskId}: ${task.action}")
executeTask(task)
if (task.recurring)
scheduleTask(client, task, task.intervalSec)
msg.ack()
}
delay(1000)
}while (true) {
auto response = client.receiveQueueMessages("scheduled-tasks", 5, 10);
for (const auto& msg : response.messages) {
auto task = ScheduledTask::fromJson(msg.body);
std::cout << "Executing task " << task.taskId << ": " << task.action << std::endl;
executeTask(task);
if (task.recurring)
scheduleTask(client, task, task.intervalSec);
msg.ack();
}
std::this_thread::sleep_for(std::chrono::seconds(1));
}use std::time::Duration;
// The unary receive auto-acknowledges due tasks as they are pulled.
loop {
let messages = client
.receive_queue_messages("scheduled-tasks", 5, 10, false)
.await?;
for m in &messages {
let task: ScheduledTask = serde_json::from_slice(&m.body).unwrap();
println!("Executing task {}: {}", task.task_id, task.action);
execute_task(&task);
if task.recurring {
schedule_task(&client, &task, task.interval_sec).await?;
println!("Task {} rescheduled in {}s", task.task_id, task.interval_sec);
}
}
tokio::time::sleep(Duration::from_secs(1)).await;
}require 'json'
# receive_queue_messages auto-acknowledges due tasks as they are pulled.
loop do
messages = client.receive_queue_messages(
channel: 'scheduled-tasks', max_messages: 5, wait_timeout_seconds: 10
)
messages.each do |m|
task = JSON.parse(m.body)
puts "Executing task #{task['taskId']}: #{task['action']}"
execute_task(task)
if task['recurring']
schedule_task(client, task, task['intervalSec'])
puts "Task #{task['taskId']} rescheduled in #{task['intervalSec']}s"
end
end
sleep 1
end# receive_queue_messages auto-acknowledges due tasks as they are pulled.
defp work_loop(client) do
case KubeMQ.Client.receive_queue_messages(client, "scheduled-tasks",
max_messages: 5,
wait_timeout: 10_000
) do
{:ok, result} ->
Enum.each(result.messages, fn m ->
task = Jason.decode!(m.body)
IO.puts("Executing task #{task["taskId"]}: #{task["action"]}")
execute_task(task)
if task["recurring"] do
schedule_task(client, task, task["intervalSec"])
IO.puts("Task #{task["taskId"]} rescheduled in #{task["intervalSec"]}s")
end
end)
{:error, err} ->
IO.puts("Receive failed: #{err.message}")
end
Process.sleep(1_000)
work_loop(client)
endAdvanced Configuration
Next Steps
Was this page helpful?