skip to content

Background Jobs & Scheduling

Getting work off the request path: task queues, background workers and scheduled jobs in .NET and Python. Asked because retries and at-least-once delivery are where it silently fails.

on this pageshow

explore

questions

page 1 of 2

In Celery, what does the celery beat process do, and how do you declare a periodic task in the beat_schedule setting?

level: juniorimportance: must knowfreq 55%

answer

  1. a clock that only sends
  2. workers still do the work
  3. entry name, task, schedule
  4. seconds, timedelta, crontab, solar
  5. celerybeat-schedule file

basics

~20 s

celery beat is a scheduler process that sends a task message whenever a beat_schedule entry comes due; workers run it. Each entry names a registered task, a schedule (seconds, timedelta, crontab or solar) and optional args, kwargs and options.

solid answer

~40 s

`celery beat` is a separate process that runs no task code. It reads the schedule, and when an entry is due it sends an ordinary task message to the broker, just as `apply_async()` would; whichever worker consumes that queue executes it. By default the entries come from `app.conf.beat_schedule`, a dict keyed by a unique entry name. Each value holds `task` (the registered task name, as a string), `schedule` (a number of seconds, a `timedelta`, a `crontab(...)` or a `solar(...)`), and optionally `args`, `kwargs` and `options`, which accepts any `apply_async()` option such as `queue` or `expires`. You start it with `celery -A proj beat`. The default `PersistentScheduler` records last-send times in a local `celerybeat-schedule` file, and only one beat may run per schedule.

code

python · 19 lines
python
from datetime import timedelta

from celery import Celery
from celery.schedules import crontab

app = Celery("saas", broker="redis://localhost:6379/0")

app.conf.beat_schedule = {
    "nightly-billing": {
        "task": "billing.tasks.run_nightly_billing",
        "schedule": crontab(hour=2, minute=0),
    },
    "hourly-cache-warmup": {
        "task": "cache.tasks.warm_cache",
        "schedule": timedelta(hours=1),
        "kwargs": {"scope": "pricing"},
        "options": {"queue": "maintenance", "expires": 50 * 60},
    },
}

go deeper

for a junior

Recall that beat only sends messages on a timetable and workers do the work. Be able to write one beat_schedule entry with a task name and a crontab or timedelta.

for a middle

Explain the entry fields, including options passing apply_async arguments, the crontab wildcard defaults, and how a timedelta schedule counts from the last send.

for a senior

Show you know beat records sends, not outcomes, so overlaps and queued backlogs are yours to handle, and that exactly one beat must run per schedule.

for a principal

Weigh a static beat_schedule in code against a database-backed schedule that operators edit, and who owns correctness when schedules change at runtime.

## What beat is, and what it is not **`celery beat`** is Celery's periodic scheduler. It is a long-running process of its own, started with `celery -A proj beat`, and its whole job is to decide *when* a message should be sent. It does **not** execute task bodies. When an entry comes due, beat publishes an ordinary task message to the broker, and from that moment the message is indistinguishable from one sent by `apply_async()` in a web request: it lands in a queue, a **worker** consuming that queue picks it up, and the worker runs the function. That split explains most beat behaviour: - If no worker is running, beat still sends on time and the messages wait in the broker. - Beat does not wait for a run to finish before sending the next one, so a slow task can overlap with its own next run. - Beat records the time it **sent** an entry (`last_run_at`), not whether the task succeeded. In the reserved scenario, a SaaS product, beat sends the nightly billing run, the hourly cache warm-up and the customers' report jobs; the workers in the application's containers do the actual billing, warming and reporting. ## Declaring entries in `beat_schedule` The default source of entries is the **`beat_schedule`** setting: a dict whose keys are unique entry names and whose values describe one periodic send each. | Field | Required | Meaning | |---|---|---| | `task` | yes | The registered task name as a string, such as `'billing.tasks.run_nightly_billing'` | | `schedule` | yes | A number of seconds, a `datetime.timedelta`, a `crontab(...)` or a `solar(...)` | | `args` | no | Positional arguments, a list or tuple | | `kwargs` | no | Keyword arguments, a dict | | `options` | no | Any `apply_async()` option: `queue`, `routing_key`, `expires`, `priority` and so on | | `relative` | no | For `timedelta` schedules, round the period to the clock instead of counting from beat's start | Two small traps live in this table. `task` is the task's **registered name**, not the function object; by default that name is built from the module path, but it is a name all the same. And a one-item `args` tuple needs its trailing comma, `(42,)`, because `(42)` is just the integer 42. ## The schedule types | Type | Example | Fires | |---|---|---| | Number or `timedelta` | `timedelta(hours=1)` | One interval after the last send; a new entry is first sent one interval after beat starts | | `crontab` | `crontab(hour=2, minute=0)` | At matching wall-clock times in the app's `timezone` (UTC by default) | | `solar` | `solar('sunset', -37.81753, 144.96715)` | At a sun event for a latitude and longitude, computed in UTC | `crontab` fields default to `'*'`, which surprises people: `crontab(day_of_week='sunday')` fires **every minute** on Sundays, because minute and hour were left as wildcards. Write `crontab(minute=0, hour=2)` when you mean 02:00 every day. Entries can also be added in code with `app.add_periodic_task(schedule, signature, name=...)`, usually from an `on_after_configure` or `on_after_finalize` signal handler. It writes into `beat_schedule` behind the scenes, and two entries built from the same signature need distinct `name`s or the second replaces the first. ## Running beat and where it keeps state 1. Start one beat process: `celery -A proj beat -l INFO`. 2. Start workers separately: `celery -A proj worker -l INFO`. 3. Let beat write its state file. The default scheduler, `celery.beat:PersistentScheduler`, keeps each entry's last send time in a `shelve` file named `celerybeat-schedule` (the `beat_schedule_filename` setting, or `-s` on the command line), so it needs a writable directory. The scheduler class is itself a setting, `beat_scheduler`. The common alternative is `django_celery_beat.schedulers:DatabaseScheduler`, which reads entries from database tables that an admin screen can edit. ## Only one beat per schedule Celery's documentation is explicit that you must make sure **only a single scheduler runs for a schedule at a time**, or tasks are sent more than once. Beat processes do not coordinate with each other and Celery ships no lock between them. That is why the worker's `-B` flag, which embeds a beat in each worker, is fine on a single node and dangerous on a fleet. ## Mistakes that show up in interviews - Saying beat "runs" the task: it only sends it. - Expecting beat to skip a send because the previous run is still going. - Leaving `crontab` fields as wildcards by accident. - Starting beat in every container "for redundancy".

  • Does Celery beat wait for the previous run of a periodic task to finish before sending the next one?
    No. Beat only sends messages on its timetable and never hears back from the workers. If the hourly warm-up takes ninety minutes, the next message is sent on time and a second worker can start it while the first is still running. When overlap matters, the task itself has to guard against it, for example with a lock it takes at the start.
  • Why does `crontab(day_of_week='sunday')` in a Celery beat entry fire far more often than once a week?
    Every `crontab` field you leave out defaults to `'*'`. With only `day_of_week` set, minute and hour are wildcards, so the entry matches every minute of every Sunday. A weekly job needs explicit fields, such as `crontab(minute=0, hour=3, day_of_week='sunday')`.
  • What does the `options` key of a Celery beat entry accept?
    Any keyword that `apply_async()` accepts: `queue`, `routing_key`, `exchange`, `priority`, `expires` and so on. Beat passes them through when it sends the message, so an entry can go to a dedicated queue or be discarded by the worker if it is not started before its expiry.

Beat is a school bell on a timetable: it rings at the scheduled time and has no idea whether the teachers are in the room. The teaching is done by the teachers, as the task is done by workers.

saying these in an interview costs you the question

  • celery beat executes the periodic task code itself
  • Every worker automatically runs its own copy of the beat schedule
  • crontab(day_of_week='sunday') runs once each Sunday
  • An entry's task field takes the decorated function object
  • Beat waits for the previous run to finish before sending again
  • Running beat in every container gives redundancy
open as a page

In Celery, what is the difference between the broker and the result backend, and does every app need a result backend?

level: juniorimportance: must knowfreq 62%

basics

~20 s

Celery's broker carries task messages from the caller to the workers; the result backend stores each task's state and return value. The broker is required; the result backend is optional and unset by default, so fire-and-forget tasks need none.

open as a page

In Celery, how does `self.retry()` re-run a failing task, and what do its `countdown`, `max_retries` and `exc` arguments control?

level: juniorimportance: must knowfreq 62%

basics

~20 s

A bound Celery task calls raise self.retry(exc=exc, countdown=N) in its except block: Celery publishes a new message for the same task id, delayed N seconds (180 by default), up to max_retries (3), then re-raises exc as the failure.

open as a page

In Celery, how does calling a task with delay() differ from apply_async(), and what does either call return to the caller?

level: juniorimportance: must knowfreq 65%

basics

~20 s

delay(*args, **kwargs) is a shortcut for apply_async(args, kwargs) with no execution options. Both serialize the arguments, publish one message to the broker and return an AsyncResult with the task id at once; neither waits for a worker.

open as a page

In Celery, how do chain, group and chord differ, and which fits resizing thumbnails in parallel and then running one step after all finish?

level: middleimportance: must knowfreq 45%

basics

~20 s

In Celery, a chain runs tasks in sequence and passes each result on; a group runs tasks in parallel; a chord is a group plus one body task that receives the list of all header results. Parallel thumbnails, then publish, is a chord.

open as a page

Why does calling AsyncResult.get() inside a Celery task raise RuntimeError, and how should a thumbnail workflow wait for another task instead?

level: middleimportance: must knowfreq 35%

basics

~20 s

Celery's AsyncResult.get() defaults to disable_sync_subtasks=True, so inside a worker task it raises RuntimeError: a task blocking on another can exhaust the pool and deadlock. Express the dependency with a chain, chord or link callback instead.

open as a page

A Celery PDF-export task sometimes hangs on a remote font server; how do soft_time_limit, time_limit and SoftTimeLimitExceeded stop it pinning a worker?

level: middleimportance: must knowfreq 40%

basics

~20 s

Celery sets no time limit by default, so a hung task holds its worker slot forever. soft_time_limit raises SoftTimeLimitExceeded inside the task so it can clean up; time_limit kills and replaces the pool process and fails the task with TimeLimitExceeded.

open as a page

In Celery, how do `autoretry_for`, `retry_backoff`, `retry_backoff_max` and `retry_jitter` combine, and what retry delays do they actually produce?

level: middleimportance: must knowfreq 48%

basics

~20 s

autoretry_for makes Celery call retry() when a listed exception escapes the task. retry_backoff sets the delay to factor x 2^retries seconds, capped by retry_backoff_max (600), and retry_jitter, on by default, draws a random delay between zero and the computed value.

open as a page

In Celery, how do you send slow video-transcoding tasks to their own queue and pin a dedicated worker to consume only that queue?

level: middleimportance: must knowfreq 50%

basics

~20 s

Map the transcode task to a queue such as transcode in task_routes, or pass queue= to apply_async, then start a worker with -Q transcode. Unrouted tasks keep going to the default queue, named celery, which other workers drain.

open as a page

In Celery, which worker pool suits CPU-heavy report rendering and which suits thousands of outbound webhook calls, and why?

level: middleimportance: must knowfreq 50%

basics

~20 s

Run CPU-heavy rendering on Celery's default prefork pool: separate processes each run their own interpreter and can keep a core busy. Run I/O-bound webhook calls on gevent or eventlet, where hundreds of greenlets wait on sockets cheaply inside one process.

open as a page

A SaaS runs `celery -A proj worker -B` in each of three identical containers and bills every customer three times nightly; why, and how should Celery beat be deployed?

level: seniorimportance: must knowfreq 45%

basics

~20 s

The -B flag embeds a beat scheduler in every worker, so three containers run three independent schedulers and each sends the billing task at 02:00. Celery has no lock between beat processes: run exactly one standalone celery beat.

open as a page

A payouts team on Celery's Redis broker raises visibility_timeout to 12 hours so delayed payouts stop running twice; what does that fix, and what does it cost?

level: seniorimportance: must knowfreq 40%

basics

~20 s

On Celery's Redis or SQS broker, a message left unacknowledged past visibility_timeout is redelivered, so long countdowns run twice. Twelve hours stops that, but a killed worker's tasks then wait up to 12 hours for redelivery.

open as a page

A Celery task that charges a customer's card sets `acks_late=True`; why can the charge now run twice, and what do `reject_on_worker_lost` and `acks_on_failure_or_timeout` change?

level: seniorimportance: must knowfreq 42%

basics

~20 s

Celery acknowledges a message just before running it by default; acks_late=True moves the ack to after the task finishes, so a worker crash mid-charge redelivers the message and the charge runs again. A killed pool child is still acked unless reject_on_worker_lost=True.

open as a page

A Celery worker busy with long report tasks holds short webhook tasks while another worker sits idle; how do worker_prefetch_multiplier and -O fair explain it?

level: seniorimportance: must knowfreq 40%

basics

~20 s

Each Celery worker reserves concurrency × worker_prefetch_multiplier messages (4 per slot by default), so webhooks wait in a busy worker's buffer. -O fair, the default since 4.0, only avoids busy children inside one worker; a lower prefetch fixes it.

open as a page

In Celery, what is a task signature, and how do .s() and .si() differ when the signature runs inside a chain?

level: juniorimportance: should knowfreq 38%

basics

~20 s

A Celery signature is one task call packed as data: task name, args, kwargs and execution options. In a chain, a .s() signature gets the previous step's result prepended to its arguments; an immutable .si() signature ignores it.

open as a page

In Celery, how do you list what each worker is running right now, and what do inspect active, reserved and scheduled each show?

level: juniorimportance: should knowfreq 35%

basics

~20 s

celery -A proj inspect active asks every worker, over the broker, which tasks it is executing. inspect reserved lists tasks it prefetched but has not started, inspect scheduled the ETA or countdown tasks it holds; messages still queued appear in none.

open as a page

What does the --concurrency option of a Celery worker control, and what value does it use when you leave it unset?

level: juniorimportance: should knowfreq 45%

basics

~20 s

Celery's --concurrency (-c, setting worker_concurrency) sets how many tasks one worker runs at once: child processes under prefork, threads or greenlets under other pools. Unset, it defaults to the number of CPUs the worker detects.

open as a page

How does django-celery-beat's DatabaseScheduler let an admin screen change Celery beat schedules at runtime, and when can an edit go unnoticed for a while?

level: middleimportance: should knowfreq 33%

basics

~20 s

DatabaseScheduler builds beat's schedule from PeriodicTask rows and checks a one-row PeriodicTasks change marker, bumped by save and delete signals, about every 5 seconds. Bulk update() or raw SQL skips those signals, so beat notices only at its next full reload.

open as a page

In Celery, how do the timezone and enable_utc settings decide when a beat entry with crontab(hour=2, minute=0) fires, and what must change afterwards?

level: middleimportance: should knowfreq 35%

basics

~20 s

crontab fields are matched against the wall clock of the app's timezone; with timezone unset and enable_utc True, that is UTC, so hour=2 means 02:00 UTC. PersistentScheduler resets itself when a stored zone changes; DatabaseScheduler needs a manual reset.

open as a page

In Celery, what does AsyncResult.get() do to the caller waiting on it, and what do its timeout and propagate arguments and forget() control?

level: middleimportance: should knowfreq 45%

basics

~20 s

Celery's AsyncResult.get() blocks the caller until the task finishes, forever by default. timeout raises TimeoutError without stopping the task; propagate re-raises the task's exception in the caller; forget() deletes the stored result from the backend.

open as a page

For a Celery payouts service, how does broker_url choose between RabbitMQ, Redis and Amazon SQS, and what does each transport trade away?

level: middleimportance: should knowfreq 48%

basics

~20 s

The scheme of Celery's broker_url picks the kombu transport: amqp:// for RabbitMQ, redis:// for Redis, sqs:// for Amazon SQS. RabbitMQ is a native broker; Redis emulates acknowledgements with a visibility timeout; SQS is managed but has no events or remote control.

open as a page

In Celery, how do the Redis, database and rpc:// result backends differ in who can read a task's result, and for how long?

level: middleimportance: should knowfreq 38%

basics

~20 s

Celery's Redis and database backends store results any process can read by task id; Redis expires them after result_expires (one day), a database only when celery.backend_cleanup runs from beat. rpc:// sends results as messages, readable once and only by the sending client.

open as a page

In Celery, what are task events and the worker's -E flag, and what does Flower need from the workers to show a task's history?

level: middleimportance: should knowfreq 30%

basics

~20 s

Task events are messages a Celery worker publishes as each task is received, started, succeeds or fails. They are off by default; celery worker -E turns them on. Flower builds its task list from that stream, so without events it shows no task history.

open as a page

Why does a Celery task's `AsyncResult.state` still read `PENDING` while a worker is running it, and which states can the task report?

level: middleimportance: should knowfreq 36%

basics

~20 s

Celery writes STARTED only when task_track_started, or the task's track_started, is enabled, and it is off by default; PENDING just means the backend has no record of the id. Built-in states: PENDING, STARTED, RETRY, SUCCESS, FAILURE, REVOKED.

open as a page

In Celery, when task_routes, a task's own queue option and apply_async(queue=) disagree, which one decides where the task is sent?

level: middleimportance: should knowfreq 30%

basics

~20 s

The apply_async(queue=) argument wins, then a queue option set on the task itself, then the first router in task_routes that returns a route. Within a task_routes dict an exact task name beats a glob pattern.

open as a page

What does bind=True change about a Celery task, and what can a bound task read from self.request while it runs?

level: middleimportance: should knowfreq 40%

basics

~20 s

With bind=True, Celery passes the task object itself as the first argument, self. Through self.request a running task reads its execution context: task id, retry count, args and kwargs, worker hostname, delivery info, eta and expires.

open as a page

In Celery 5.6, which task argument types survive the default JSON serializer, and what does switching a task to pickle cost?

level: middleimportance: should knowfreq 45%

basics

~20 s

Celery's default JSON carries strings, numbers, booleans, None, lists and dicts, plus datetime, Decimal, UUID and bytes, which kombu tags and restores; anything else raises EncodeError at the call. Pickle carries more but lets broker writers run code on workers.

open as a page

Why would a Celery worker log 'Received unregistered task of type' for an invoice task, and how does Celery name tasks by default?

level: middleimportance: should knowfreq 35%

basics

~20 s

A Celery message carries only the task's name, looked up in the worker's own registry. Names default to module path plus function name, so a worker that never imported the module, or imported it under another path, rejects it as unregistered.

open as a page

How do Celery's worker_max_tasks_per_child and worker_max_memory_per_child stop report-rendering workers from growing, and what do they cost?

level: middleimportance: should knowfreq 35%

basics

~10 s

Both recycle Celery prefork child processes: worker_max_tasks_per_child replaces a child after N tasks, worker_max_memory_per_child (kilobytes) replaces it after the task that crossed the limit finishes. Each restart costs a fork and warm-up.

open as a page

showing 1–30 of 43