Skip to content

RQ

Requires RQ 1.10.1+, Python 3.11+. See the compatibility matrix for the full pin string.

Terminal window
pip install z4j-rq

RQ does not expose a worker middleware system. z4j wraps the class that owns Worker.execute_job, the parent-side execution boundary, and also provides optional per-job callback functions.

Hook z4j event
Before Worker.execute_job task.started
Finished job status after Worker.execute_job task.succeeded
Failed or stopped status after Worker.execute_job task.failed
Explicit per-job stopped callback task.revoked
Verb How
submit enqueue the import-path task name on the selected queue (Queue.enqueue_at when an ETA is given); a priority is refused, route priority classes to separate queues
retry requeue a failed job by reference; complete operator-supplied replacements use a new enqueue; refused while the job is still running; an ETA needs both replacement collections and must lie between 60 s in the past and one year ahead
cancel send_stop_job_command if running (RQ 1.13+); job.cancel() if queued, deferred, or scheduled; a job already finished, failed, canceled or stopped returns success with noop: true
purge_queue queue.empty(); 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
bulk_retry retry explicit project-owned task IDs only, capped at 10,000
requeue_dead_letter requeue from RQ's FailedJobRegistry
list_dead_letters page each queue's FailedJobRegistry newest first (job id, function name, queue, failure time, redacted traceback tail) without deserialising any job payload; the dashboard's dead-letters page and GET /projects/{slug}/dead-letters read it

RQ uses many queues per app. The adapter discovers queues through Queue.all(connection) when it builds queue and task snapshots. Each queue appears separately in the dashboard with its own counts.

rq-worker-pool works fine - each worker in the pool connects as its own worker (the runtime builds worker_id from the framework name, the PID and the start time) beneath the one agent its token identifies. Keep agent_name stable; it only labels the shared agent row.

RQ has no worker-ready signal, so the worker process calls the bootstrap itself from a module the worker imports (for example settings.py or rq_settings.py):

from z4j_rq import register_worker_bootstrap
register_worker_bootstrap()

It reads Z4J_BRAIN_URL / Z4J_TOKEN / Z4J_PROJECT_ID / Z4J_HMAC_SECRET from the environment, connects to Z4J_RQ_REDIS_URL (default redis://localhost:6379/0), and calls install_agent with an RqEngineAdapter. Z4J_DISABLED=1 skips it, a truthy Z4J_DEV_MODE is forwarded as dev_mode=True, and a second call in the same process is a no-op. It also skips, with an INFO log, unless Z4J_BRAIN_URL, Z4J_TOKEN and Z4J_PROJECT_ID are all set, and a failed Redis ping aborts the bootstrap with a warning while the worker continues without z4j.

For POSIX fork workers, z4j-rq provides a corrected worker using RQ's supported worker-class option:

Terminal window
rq worker --worker-class z4j_rq.worker.Worker default

Keep your application's queue names, Redis options, settings module and z4j startup integration. Applications that construct workers in Python can use from z4j_rq.worker import Worker with their existing constructor arguments.

Stock RQ assigns every new RQ_JOB_ID to the long-lived parent's environment as well as the child's, which grows parent memory on Linux: glibc retains distinct environment values. The corrected worker makes those assignments only in the short-lived child, where the job still receives its worker ID and job ID. It inherits RQ's job execution, monitoring, timeouts, callbacks, retries and shutdown, preserving RQ 1.x session isolation and RQ 2.x process-group isolation.

This class must be selected explicitly. Installing z4j-rq does not change stock/custom workers, and this repair does not replace SimpleWorker, SpawnWorker or Windows worker implementations. Custom classes that override the fork boundary need their own integration review.

For an affected stock fork worker that has not switched to the corrected class, the fallback is rq worker --max-jobs 50000 default under a process manager that restarts successful exits (Restart=always for systemd or autorestart=true for Supervisor). Allow current jobs to finish during shutdown. Keep normal process supervision, memory monitoring and job idempotency with either worker; fixing this allocation does not prevent memory growth caused by application code or establish unlimited worker lifetimes.

  • Failed queue - RQ moves failures to a FailedJobRegistry. z4j uses that registry to list dead letters and for the retry and dead-letter requeue actions, so a parked job can be inspected and put back on its queue from the dashboard; lifecycle visibility comes from the worker hook.
  • Scheduler actions - the rq-scheduler companion supports list, trigger, destructive disable, and delete. It does not support create, update, or enable; re-enable by registering the job again from application code.

Pass the application's queue, scheduler, or another object that exposes its Redis connection to the adapter:

from z4j_rq import RqEngineAdapter
adapter = RqEngineAdapter(rq_app=queue)

See scheduler: rq-scheduler.