Celery
Requires Celery 5.2.2+, Python 3.11+. See the compatibility matrix for the full pin string.
z4j is a Celery dashboard, monitoring, and control plane that captures task lifecycle events, persists them in Postgres with an HMAC-chained audit log, and lets operators retry, cancel, or schedule tasks from a unified UI. See the marketing landing for the Celery feature overview and the Flower alternative page if you're migrating from Flower.
Install
Section titled “Install”pip install z4j-celeryPip transitively pulls z4j-core and z4j-bare. No explicit adapter registration is needed if you use z4j-django / z4j-flask / z4j-fastapi. The worker_ready signal handler is wired:
- by
z4j-django(which eagerly importsz4j_celeryat module-load time when inINSTALLED_APPS), or - when a Flask or FastAPI application imports its z4j framework package and the Celery worker imports that application module, or
- by direct
import z4j_celeryin your Celery application module.
Celery does not discover this package through a Celery plugin entry point. Ensure the worker process imports one of the paths above.
Worker bootstrap
Section titled “Worker bootstrap”When a Celery worker is ready, the worker_ready signal fires and z4j's worker bootstrap:
- Returns without starting anything when
Z4J_DISABLEDis truthy (1,true,yes,on). - Confirms from
sys.argvthat the process is acelery workerinvocation (celery,python -m celery,uvx celery, or auv/pipx/poetry/hatch/rye/pdmrun launcher in front of one of those), notcelery inspect/control/purge/etc. - those don't deserve an agent slot. - Resolves the Celery app from
sender.apporcurrent_app. - Calls
install_agent(engines=[CeleryEngineAdapter(celery_app=...)]), which readsZ4J_BRAIN_URL/Z4J_TOKEN/Z4J_HMAC_SECRET/Z4J_PROJECT_ID/Z4J_AGENT_NAMEfrom the environment. A truthyZ4J_DEV_MODEin the worker environment is forwarded asdev_mode=True. - Logs
INFO:z4j.adapter.celery.worker_bootstrap:z4j worker bootstrap: agent runtime started (celery_app=..., framework=...)on success.framework=names the first framework adapter package found amongz4j-fastapi,z4j-flaskandz4j-django, elsebare; it only labels the agent row.
If you don't see that log line on worker boot, the agent is not running and tasks will not reach z4j. When no Z4J_* config is set, the bootstrap logs z4j worker bootstrap: no Z4J_* env config found, skipping auto-start at INFO; any other startup failure logs z4j worker bootstrap: failed to start agent runtime - worker will continue without z4j observability. The worker_ready handler is connected whenever z4j_celery is imported, but it returns silently unless the command line looks like one of the worker invocations above, so a worker started from a script with app.worker_main() gets no agent and must call install_agent() itself.
What it captures
Section titled “What it captures”| Signal | z4j event |
|---|---|
task_received |
task.received |
task_prerun |
task.started |
task_postrun (state=SUCCESS) |
task.succeeded |
task_failure |
task.failed |
task_retry |
task.retried |
task_revoked |
task.revoked |
Plus the payload: args and kwargs (redacted), the queue (taken from delivery_info.routing_key), and parent_task_id / root_task_id from Celery's canvas linkage (chord / group support).
Under prefork, gevent and eventlet pools the task signals fire inside the child processes, and the adapter drops a signal event raised in any process other than the one that connected the hooks. To cover those pools it turns on worker_send_task_events and task_send_sent_event, broadcasts control.enable_events(), and runs a broker-events consumer in the worker's main process that maps task-started / task-succeeded / task-failed / task-retried / task-revoked; under the solo pool the two sources overlap and the brain deduplicates them. Set Z4J_CELERY_LOCAL_HOSTNAME (or CELERY_HOSTNAME) to this worker's node name so the consumer ignores events from other workers on the same broker; without it every agent receives every event and the brain relies on its own dedupe.
Events wait in an in-process queue of 10,000 entries before the runtime drains them into the local buffer. When it is full the oldest entry is dropped and a warning is logged.
Actions
Section titled “Actions”| Verb | How z4j performs it |
|---|---|
retry |
app.send_task(...); reuses result-backend inputs when result_extended stored them, or requires complete operator replacements |
cancel |
fire-and-forget app.control.revoke(task_id, terminate=True, signal="SIGTERM") |
purge_queue |
channel.queue_purge(queue_name) on a write connection; refused without a confirm token (HMAC keyed with Z4J_HMAC_SECRET, bound to the queue's observed depth), when the depth cannot be measured, or when the depth exceeds Z4J_PURGE_THRESHOLD (default 1000), unless force=True |
bulk_retry |
explicit task_ids only (no broker walk), one native retry at a time, capped at 10k; after 3 consecutive broker timeouts the remaining ids are skipped and the result carries circuit_broken: true |
submit_task |
apply_async on a locally registered task, app.send_task for unknown names |
restart_worker |
app.control.broadcast("pool_restart", arguments={"reload": True}) to one worker; restarts the pool, not the parent process, and only when Celery's worker_pool_restarts is enabled |
pool_grow / pool_shrink |
app.control.pool_grow(delta) / pool_shrink(delta) addressed to one worker |
add_consumer / cancel_consumer |
app.control.add_consumer(queue) / cancel_consumer(queue) addressed to one worker |
rate_limit |
validates the rate against Celery's grammar (0 clears; <n>, <n>/s, <n>/m, <n>/h) before control.rate_limit. At the adapter API an omitted worker_name broadcasts to every worker, but a command from the dashboard or REST API never reaches that path: the agent dispatcher fills a missing worker name from the command target id, which the brain sets to the task name, so the control command goes to a worker named after the task and can report success without changing the fleet. Pass an explicit worker name |
Dead-letter listing and requeue are not advertised. Celery has no portable
dead-letter store, and a generic retry cannot consume or acknowledge a
broker-specific dead-letter entry. The brain refuses a requeue for a Celery
engine with 422 and a listing with 409, and nothing touches the broker.
Chord / group support
Section titled “Chord / group support”Celery stamps parent_id and root_id on every task request (the caller, and the entry point of the chain, group or chord). z4j forwards them as parent_task_id and root_task_id, letting the dashboard:
- Show sub-tasks grouped under the parent.
- Render the whole chain, group or chord as a tree on the task detail page.
Caveats
Section titled “Caveats”- Canvas signatures -
.si()/.s()are not reconstructed as signature objects. Retry uses complete unredacted inputs from the result backend when available, or complete operator-supplied replacements; a redacted event payload is not replay authority. - eta/countdown on retry - the original ETA is not recovered automatically; an operator-supplied retry ETA is honored when it lies between 60 s in the past and 365 days ahead, and refused (never clipped) outside that window.
- Broker state visibility - queue depth comes from
LLENon a Redis broker and from a passivequeue_declareon an AMQP broker; no management plugin is involved. Depth is probed only for the queues intask_queues, ortask_default_queuewhen that is unset, and worker stats come fromcontrol.inspectat most once every 60 s.
Config
Section titled “Config”The adapter has no settings.Z4J keys of its own. It reads the Celery app's config directly (worker_pool, broker_url, task_queues / task_default_queue) and these environment variables: Z4J_CELERY_LOCAL_HOSTNAME or CELERY_HOSTNAME (filter the broker-events consumer to this worker), Z4J_PURGE_THRESHOLD (purge refusal depth, default 1000), Z4J_HMAC_SECRET (keys the purge confirm token) and Z4J_ACCEPT_LEGACY_PURGE_TOKEN (accept the unkeyed purge token from an older brain during a rolling upgrade; default off).
For Django, the app is auto-detected from settings.CELERY_APP, the project package's celery_app / app attribute, celery.current_app, and <package>.celery.app (see quickstart §Auto-detect). If your layout is unusual:
# settings.py - top-level Django setting, NOT inside the Z4J dictCELERY_APP = "myproject.celery:app"Flask resolves the app from app.config["CELERY_APP"] (an object or an import path), app.extensions["celery"] or an app.celery_app attribute; FastAPI takes celery_app= on z4j_lifespan() / install_z4j(). For bare-Python layouts, pass the app directly to the engine:
from z4j_celery.engine import CeleryEngineAdapterfrom z4j_bare.install import install_agent
install_agent( engines=[CeleryEngineAdapter(celery_app=my_celery_app)], # brain_url / token / hmac_secret / project_id / agent_name # are read from Z4J_* env vars if not passed explicitly)Performance notes
Section titled “Performance notes”- Signal handlers only map and queue events locally; they make no brain network request inline.