# Scheduled Task System (/learn/queues/scenarios/scheduled-tasks)



## Architecture [#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.

<Mermaid
  chart="graph LR
  SCHED[&#x22;Task Scheduler&#x22;]
  DQ{{&#x22;delayed buffer<br/>delay=60s&#x22;}}
  TQ{{&#x22;tasks queue&#x22;}}
  TW[&#x22;Task Worker&#x22;]
  DONE[&#x22;Complete&#x22;]

  SCHED -- &#x22;send (delaySeconds)&#x22; --> DQ
  DQ -- &#x22;delay expires&#x22; --> TQ
  TQ -- &#x22;receive + process&#x22; --> TW
  TW -. &#x22;recurring? re-enqueue&#x22; .-> SCHED
  TW -- &#x22;one-shot&#x22; --> DONE

  class DQ,TQ queue
  class SCHED,TW,DONE client"
/>

*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 [#schedule-a-task]

Send a message with `delaySeconds` set to the desired future execution time.

<Tabs groupId="language" items="['Go', 'Python', 'Node.js', 'Java', 'C#', 'Kotlin', 'C++', 'Rust', 'Ruby', 'Elixir']">
  <Tab value="Go">
    ```go title="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
    }
    ```
  </Tab>

  <Tab value="Python">
    ```python title="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)
    ```
  </Tab>

  <Tab value="Node.js">
    ```typescript title="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}`);
    }
    ```
  </Tab>

  <Tab value="Java">
    ```java title="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());
    }
    ```
  </Tab>

  <Tab value="C#">
    ```csharp title="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}");
    }
    ```
  </Tab>

  <Tab value="Kotlin">
    ```kotlin title="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}")
    }
    ```
  </Tab>

  <Tab value="C++">
    ```cpp title="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;
    }
    ```
  </Tab>

  <Tab value="Rust">
    ```rust title="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(())
    }
    ```
  </Tab>

  <Tab value="Ruby">
    ```ruby title="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)
    ```
  </Tab>

  <Tab value="Elixir">
    ```elixir title="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
    ```
  </Tab>
</Tabs>

## Task Worker with Recurring Support [#task-worker-with-recurring-support]

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

<Tabs groupId="language" items="['Go', 'Python', 'Node.js', 'Java', 'C#', 'Kotlin', 'C++', 'Rust', 'Ruby', 'Elixir']">
  <Tab value="Go">
    ```go title="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()
    }
    ```
  </Tab>

  <Tab value="Python">
    ```python title="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)
    ```
  </Tab>

  <Tab value="Node.js">
    ```typescript title="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));
    }
    ```
  </Tab>

  <Tab value="Java">
    ```java title="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);
    }
    ```
  </Tab>

  <Tab value="C#">
    ```csharp title="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);
    }
    ```
  </Tab>

  <Tab value="Kotlin">
    ```kotlin title="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)
    }
    ```
  </Tab>

  <Tab value="C++">
    ```cpp title="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));
    }
    ```
  </Tab>

  <Tab value="Rust">
    ```rust title="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;
    }
    ```
  </Tab>

  <Tab value="Ruby">
    ```ruby title="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
    ```
  </Tab>

  <Tab value="Elixir">
    ```elixir title="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
    ```
  </Tab>
</Tabs>

## Advanced Configuration [#advanced-configuration]

<Accordions>
  <Accordion title="Maximum Delay">
    The server setting `MaxDelaySeconds` (default: 43,200 = 12 hours) limits the maximum delay. For longer intervals, re-enqueue with the max delay and decrement a remaining counter.
  </Accordion>

  <Accordion title="Delay + Expiration">
    Combine delay with expiration for tasks that should only run within a window. If the delay is 1 hour and expiration is 2 hours, the task must be processed within 1 hour of becoming available.
  </Accordion>
</Accordions>

## Next Steps [#next-steps]

<Cards>
  <Card title="Delayed Messages" href="/learn/queues/tutorials/delayed-messages" description="Core delayed message mechanics." />

  <Card title="Webhook Delivery" href="/learn/queues/scenarios/webhook-delivery" description="Reliable delivery with retry and backoff." />
</Cards>
