Draining RQ Workers Safely

RQ is deliberately simple, and its shutdown behaviour is simple too — which means the gaps are yours to close. This guide shows what RQ does on SIGTERM and SIGKILL, where interrupted jobs end up, and how to deploy RQ workers without losing work, as part of Graceful Shutdown & Worker Deployments in Backend Frameworks & Worker Scaling.

Problem Statement

A Flask application runs RQ workers in Kubernetes. After deploys, some jobs disappear: they are not in the queue, not in the finished registry, and not in the failed registry — only visible as "started" long after they should have ended. Others are marked failed with Work-horse terminated unexpectedly and need manual requeueing. The team wants deploys that let running jobs finish, a known outcome for every interrupted job, and automatic recovery rather than someone checking registries after each deploy.

Prerequisites

  • RQ 1.15+ or 2.x with Redis.
  • Access to the worker's Deployment spec and entrypoint.
  • Job timeouts (job_timeout) set for each job type.
  • Idempotent jobs, because recovered jobs run again.

Step 1 — Understand the Fork Model

An RQ worker is a parent process that dequeues a job and forks a work horse child to run it. The job is moved to the queue's StartedJobRegistry while the horse runs. On success it goes to FinishedJobRegistry; on exception, to FailedJobRegistry.

worker (parent) --dequeue--> job moves to StartedJobRegistry
        |--fork--> work horse runs job --> Finished / Failed registry

Shutdown signals go to the parent:

  • First SIGTERM or SIGINT (warm shutdown): the parent stops dequeuing and waits for the current horse to finish, then exits.
  • Second signal (cold shutdown): the parent kills the horse immediately and exits; the job is marked failed (or left in the started registry if the parent itself dies).
Parent, horse, and two kinds of shutdown The RQ parent process dequeues a job, records it in the started registry, and forks a work horse that runs it. On the first SIGTERM, the parent stops dequeuing and waits for the horse to finish. On a second signal, or a SIGKILL of the parent, the horse is killed; the job either lands in the failed registry or remains stranded in the started registry. RQ worker processes parent dequeue, fork, signals work horse runs one job 1st TERM: wait for horse (warm) 2nd signal / KILL: horse dies Kubernetes sends one TERM, then SIGKILL at the end of the grace period.

Step 2 — Give Warm Shutdown Enough Time

Kubernetes sends one SIGTERM and, after the grace period, SIGKILL. That maps to RQ's warm shutdown followed by an uncatchable kill. The grace period must exceed the longest job that can be running — which, for RQ, is bounded by job_timeout.

spec:
  terminationGracePeriodSeconds: 330      # > longest job_timeout (300) + margin
  containers:
    - name: rq-worker
      command: ["rq", "worker", "--url", "$(REDIS_URL)", "high", "default",
                "--with-scheduler", "--logging_level", "INFO"]
# Every enqueue sets a timeout that fits inside the grace period
q.enqueue(generate_invoice, invoice_id, job_timeout=120)
q.enqueue(rebuild_search_index, job_timeout=300)

Jobs with job_timeout longer than the grace period will be SIGKILLed on every deploy that catches them mid-run. Either split them, or run them on a separate worker deployment with a longer grace period.

Step 3 — Know Where Killed Jobs End Up

When the pod is SIGKILLed, the parent cannot record anything. The job stays in StartedJobRegistry — "started" forever from RQ's point of view. RQ's registry cleanup eventually notices started jobs whose worker heartbeat has expired and moves them to the failed registry (AbandonedJobError in RQ 2.x), but only when cleanup runs.

from rq import Queue
from rq.registry import StartedJobRegistry, FailedJobRegistry

q = Queue("default", connection=redis)
started = StartedJobRegistry(queue=q)
print("started:", started.get_job_ids())
started.cleanup()                                  # moves expired started jobs to failed
failed = FailedJobRegistry(queue=q)
for job_id in failed.get_job_ids():
    job = q.fetch_job(job_id)
    print(job_id, job.exc_info.splitlines()[-1] if job.exc_info else "")

Neither outcome retries the job automatically. That is the "disappearing jobs" symptom: they are stranded in a registry nobody checks.

Three outcomes for a running job If the job finishes within the grace period, it reaches the finished registry. If a cold shutdown kills the horse while the parent survives, the job goes to the failed registry with a work-horse error. If the whole pod is SIGKILLed, the job stays in the started registry until registry cleanup moves it to failed as abandoned. None of the failure outcomes retry automatically. Where an interrupted job lands warm, finished in time FinishedJobRegistry cold shutdown FailedJobRegistry no automatic retry pod SIGKILLed stranded in Started until cleanup: failed Step 4 turns the two right-hand outcomes into automatic retries.

Step 4 — Retry Interrupted Jobs Automatically

Give jobs a Retry policy so ordinary failures retry, and run a small reaper that requeues jobs that were abandoned or killed by shutdown.

from rq import Retry

q.enqueue(generate_invoice, invoice_id, job_timeout=120,
          retry=Retry(max=3, interval=[30, 120, 600]))

# reaper.py — run every minute (rq-scheduler, cron, or a periodic job)
INTERRUPTION_MARKERS = ("AbandonedJobError", "Work-horse terminated unexpectedly",
                        "Work-horse was terminated unexpectedly")

def requeue_interrupted(queue: Queue, max_requeues: int = 3):
    StartedJobRegistry(queue=queue).cleanup()               # move abandoned started jobs to failed
    failed = FailedJobRegistry(queue=queue)
    for job_id in failed.get_job_ids():
        job = queue.fetch_job(job_id)
        if not job or not job.exc_info or not any(m in job.exc_info for m in INTERRUPTION_MARKERS):
            continue                                        # real failures stay for triage
        count = int(job.meta.get("interrupt_requeues", 0))
        if count >= max_requeues:
            continue                                        # repeatedly killed: needs a human
        job.meta["interrupt_requeues"] = count + 1
        job.save_meta()
        failed.requeue(job_id)

Distinguishing interruptions from genuine failures keeps real bugs visible in the failed registry while deploy-caused interruptions heal themselves. The count limit prevents a job that kills its horse every time (for example, by exhausting memory) from looping forever.

Step 5 — Separate Long Jobs and Keep the Parent Healthy

Put long jobs on their own queues served by their own worker deployment with a longer grace period, so normal deploys of short-job workers stay fast.

# rq-worker-long: only the "long" queue, 30-minute grace period
command: ["rq", "worker", "long", "--url", "$(REDIS_URL)"]
terminationGracePeriodSeconds: 1860      # job_timeout 1800 + margin

Also keep the parent process's heartbeat healthy: worker heartbeats are how RQ decides a started job is abandoned. The default --worker-ttl (420 seconds) means cleanup waits that long before treating a dead worker's job as abandoned; lower it if you want faster recovery. Scaling RQ fleets is covered in scaling RQ workers in production.

Step 6 — Make the Entrypoint Forward Signals

The most common reason RQ never sees SIGTERM is an entrypoint shell script that is PID 1 and does not forward signals.

#!/bin/sh
# entrypoint.sh
set -e
python manage.py wait_for_redis
exec rq worker --url "$REDIS_URL" high default     # exec: rq becomes PID 1 and receives TERM

Without exec, the shell receives SIGTERM, ignores it, and Kubernetes SIGKILLs everything after the grace period — every running job stranded on every deploy.

Step 7 — Watch the Registries

Registries are RQ's only record of jobs that did not finish, so they belong on a dashboard. Export their sizes per queue and alert when the started registry holds jobs older than the longest job_timeout, or when the failed registry grows.

from prometheus_client import Gauge
STARTED = Gauge("rq_started_jobs", "", ["queue"])
FAILED = Gauge("rq_failed_jobs", "", ["queue"])
OLDEST_STARTED = Gauge("rq_oldest_started_seconds", "", ["queue"])

def export_registry_metrics(queue):
    reg = StartedJobRegistry(queue=queue)
    STARTED.labels(queue.name).set(reg.count)
    FAILED.labels(queue.name).set(FailedJobRegistry(queue=queue).count)
    oldest = min((j.started_at for j in map(queue.fetch_job, reg.get_job_ids()) if j and j.started_at),
                 default=None)
    OLDEST_STARTED.labels(queue.name).set((utcnow() - oldest).total_seconds() if oldest else 0)
Self-healing after deploys Every minute the reaper runs registry cleanup, which moves started jobs whose worker heartbeat expired into the failed registry. It then scans failed jobs: those whose error indicates interruption by shutdown are requeued, up to three times per job. Jobs that failed for real reasons, or were interrupted too many times, stay in the failed registry and raise an alert. Reaper, every minute cleanup() abandoned to failed interruption? requeue (max 3) real failure stays, alerts Deploy interruptions heal automatically; genuine bugs remain visible for triage.

Verification

# Start a long job, trigger a rollout, then check registries
python -c 'from app import q, slow_job; q.enqueue(slow_job, 90, job_timeout=120)'
kubectl rollout restart deployment/rq-worker
sleep 150
rq info --url "$REDIS_URL"                    # started registry should be empty
python -c 'from app import q; from rq.registry import FinishedJobRegistry; print(FinishedJobRegistry(queue=q).count)'

The job should appear in the finished registry (warm shutdown let it complete) or be requeued by the reaper — never remain in started.

Gotchas & Edge Cases

job_timeout longer than the grace period. Guarantees SIGKILL of long jobs on deploy. Keep the two aligned per worker deployment.

Windows and SimpleWorker. SimpleWorker runs jobs in-process without forking; a crash takes the worker with it and there is no horse to kill separately. Use the default forking worker on Linux.

Dependent jobs. Jobs enqueued with depends_on wait in the deferred registry until their parent finishes; if the parent is stranded in started, dependents wait too. The reaper's requeue unblocks the chain once the parent completes.

Registry growth. Finished and failed registries have TTLs (result_ttl, failure_ttl); set them so registries do not grow without bound.

Scheduler. With --with-scheduler, one worker also moves scheduled jobs to queues; it takes a lock, so multiple workers with the flag are safe.

FAQ

Is RQ's lack of automatic recovery a reason to switch to Celery? Not necessarily; a reaper like the one above covers it. If you need many features like this, compare frameworks as in RQ vs Celery for Python.

Can a job detect shutdown and stop early? The horse receives no warning in warm shutdown; it simply finishes. For long jobs, checkpoint progress so an eventual interruption is cheap.

How do I drain a worker without deploying? Send SIGTERM to the worker's parent process (or rq suspend to pause every worker via a Redis flag, then rq resume). rq suspend is useful for broker maintenance: workers finish their current jobs and stop dequeuing until resumed, and jobs simply accumulate in the queues meanwhile.

Should the reaper run inside a worker? Any single process works — an rq-scheduler periodic job, a cron container, or a small sidecar — as long as only one instance runs at a time. failed.requeue is safe to call twice for the same job, but running the reaper everywhere wastes Redis work and makes the requeue counts noisy.

What about jobs enqueued with at_front=True? Requeued jobs go to the back of their queue by default. If interrupted jobs should run before newer ones, requeue them with at_front=True, keeping in mind that a large burst of requeues can then delay fresh work.

Why do started jobs stay "started" so long? Because cleanup waits for the worker's heartbeat TTL to expire. Lower --worker-ttl for quicker detection.

Related