KubeMQ
LearnQueuesScenarios

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.

scheduler.go
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
}
scheduler.py
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)
scheduler.ts
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}`);
}
Scheduler.java
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());
}
Scheduler.cs
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}");
}
Scheduler.kt
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}")
}
scheduler.cpp
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;
}
scheduler.rs
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(())
}
scheduler.rb
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)
scheduler.exs
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}")
end

Task Worker with Recurring Support

When a task completes, check if it's recurring and re-schedule it.

task_worker.go
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()
}
task_worker.py
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)
task_worker.ts
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));
}
TaskWorker.java
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);
}
TaskWorker.cs
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);
}
TaskWorker.kt
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)
}
task_worker.cpp
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));
}
task_worker.rs
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;
}
task_worker.rb
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
task_worker.exs
# 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)
end

Advanced Configuration

Next Steps

Was this page helpful?

On this page