Skip to content

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.

Terminal window
pip install z4j-celery

Pip 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 imports z4j_celery at module-load time when in INSTALLED_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_celery in 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.

When a Celery worker is ready, the worker_ready signal fires and z4j's worker bootstrap:

  1. Returns without starting anything when Z4J_DISABLED is truthy (1, true, yes, on).
  2. Confirms from sys.argv that the process is a celery worker invocation (celery, python -m celery, uvx celery, or a uv / pipx / poetry / hatch / rye / pdm run launcher in front of one of those), not celery inspect/control/purge/etc. - those don't deserve an agent slot.
  3. Resolves the Celery app from sender.app or current_app.
  4. Calls install_agent(engines=[CeleryEngineAdapter(celery_app=...)]), which reads Z4J_BRAIN_URL / Z4J_TOKEN / Z4J_HMAC_SECRET / Z4J_PROJECT_ID / Z4J_AGENT_NAME from the environment. A truthy Z4J_DEV_MODE in the worker environment is forwarded as dev_mode=True.
  5. 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 among z4j-fastapi, z4j-flask and z4j-django, else bare; 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.

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.

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.

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.
  • 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 LLEN on a Redis broker and from a passive queue_declare on an AMQP broker; no management plugin is involved. Depth is probed only for the queues in task_queues, or task_default_queue when that is unset, and worker stats come from control.inspect at most once every 60 s.

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 dict
CELERY_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 CeleryEngineAdapter
from 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
)
  • Signal handlers only map and queue events locally; they make no brain network request inline.

See scheduler: celery-beat.