Skip to content

Dramatiq

Requires Dramatiq 1.14+, Python 3.11+. See the compatibility matrix for the full pin string.

Terminal window
pip install z4j-dramatiq

z4j registers a Middleware on the broker. Do not call connect_signals() yourself: install_agent (and the Flask and FastAPI adapters) call it once at start, and a second call installs a second middleware and duplicates every event.

from dramatiq import get_broker
from z4j_bare import install_agent
from z4j_dramatiq import DramatiqEngineAdapter
broker = get_broker()
runtime = install_agent(engines=[DramatiqEngineAdapter(broker=broker)])

The Flask adapter (app.config["DRAMATIQ_BROKER"]) and the FastAPI adapter (z4j_lifespan(dramatiq_broker=...)) register the adapter automatically, and both also adopt dramatiq.get_broker() when at least one actor is registered on it. z4j-django discovers only Celery, so Django and bare deployments pass the adapter to install_agent themselves (see Config below).

The adapter reports no queue or worker list (list_queues and list_workers return empty lists). Queue depths reach the dashboard only through the heartbeat's health payload, which reads broker.get_queue_message_counts for each queue that has an actor registered in this process.

Middleware hook z4j event
after_enqueue task.received
before_process_message task.started
after_process_message (no exception) task.succeeded
retry chosen by Retries task.retried
exhausted exception task.failed
Verb How
submit sends through the registered actor, or a raw message for a remote actor
retry re-sends only complete operator-supplied replacement args and kwargs
cancel prevents pending work from starting, but only with z4j-dramatiq[abort] and a real dramatiq_abort.Abortable middleware in this broker's stack
purge_queue broker.flush(queue_name) (Dramatiq 1.16+); refused without a confirm token (HMAC keyed with Z4J_HMAC_SECRET, bound to the observed depth), when the depth cannot be measured, or above Z4J_PURGE_THRESHOLD (default 10,000), unless force; a flush that exceeds its 10 s timeout reports an indeterminate result
list_dead_letters read each queue's <queue>.XQ dead-letter store without consuming it; promoted on the Redis, RabbitMQ and stub brokers. Redis pages full entries (message id, actor, dead-letter time, attempts, redacted traceback tail); RabbitMQ reports the count only, with an empty page, because AMQP has no non-destructive read of a queue

Stock Dramatiq does not advertise cancel. The adapter promotes it only when the abort extra is installed and the broker contains a real dramatiq_abort.Abortable instance. Configure that middleware using dramatiq-abort's event backend before constructing the adapter. z4j always requests its pending-only AbortMode.CANCEL: a pending message is prevented from starting, but running work is never interrupted. Without both the package and middleware, a direct cancel call fails closed. A cancel that exceeds its 10 s broker timeout reports an indeterminate result, since the pending-task cancellation may still have landed.

Dead letters can be listed (the dashboard's dead-letters page and GET /projects/{slug}/dead-letters read the .XQ store), but bulk retry and dead-letter requeue remain unavailable because Dramatiq has no portable primitive that satisfies those z4j contracts: there is no resurrect-by-id API, so the adapter does not advertise requeue_dead_letter and a requeue request is refused before anything touches the broker. Message bodies are decoded as JSON only; a deployment using Dramatiq's PickleEncoder lists id, queue and time with an empty actor name and error excerpt, because the agent never unpickles broker data.

install_agent calls DramatiqEngineAdapter.connect_signals() once at start. It inserts Z4JMiddleware at the head of the broker's middleware list, before Retries, because Dramatiq runs after_* hooks in reverse order. This lets Retries make its decision before z4j classifies the attempt. Do not call connect_signals() yourself: each call adds a fresh middleware (Dramatiq allows duplicates) and every event is then emitted twice.

  • No chord/group primitive in Dramatiq, so there is no chord-aware UI.
  • z4j rejects submit or retry requests carrying eta or priority; the adapter cannot honor those options portably across its supported brokers.
  • Submit and retry requests naming a queue other than the actor's registered queue are rejected; Dramatiq cannot override an actor's queue at send time. That check covers only actors registered in this process; for an unregistered actor name the adapter enqueues a raw Message on the named queue, or default when none is given.

The adapter takes the broker object; there is no settings.Z4J key for it. An unknown key under Z4J fails config validation, so z4j-django logs failed to build agent runtime and starts no agent.

# Flask: the extension reads app.config["DRAMATIQ_BROKER"]
# (a broker object or an import path)
app.config["DRAMATIQ_BROKER"] = "myapp.broker.redis_broker"
# FastAPI: hand the broker to the lifespan
app = FastAPI(lifespan=z4j_lifespan(dramatiq_broker=broker))
# Django and bare Python: register the adapter yourself in the
# Dramatiq worker process
from z4j_bare import install_agent
from z4j_dramatiq import DramatiqEngineAdapter
runtime = install_agent(engines=[DramatiqEngineAdapter(broker=broker)])

See scheduler: APScheduler - the typical scheduler for Dramatiq apps.