07 · Long-Running Work: Job Queues and Polling¶
Some work doesn't fit in a request: generating a big report, transcoding a video, sending
a thousand emails, calling a slow partner API. Level 3 lesson 5 showed that
BackgroundTasks runs after the response but in the web process, with no
persistence, no retries and no visibility. When the work must happen, it needs a
job queue: the API records the job durably and answers immediately; separate
worker processes do the work; the client polls (or gets notified) for the result.
Production systems usually use a queue library with a broker — Celery, RQ or Dramatiq (Redis or RabbitMQ), ARQ (Redis, asyncio), or a managed cloud queue. No broker was available for this course, so this lesson builds the essential mechanics with a SQLite table. That's genuinely useful for small systems (a PostgreSQL version of the same table is a well-known pattern), and it makes every moving part visible — the same parts you configure in Celery or ARQ.
The protocol¶
POST /reports→ 202 Accepted, a job ID, andLocation: /jobs/{id}.GET /jobs/{id}→queued/running(withRetry-After) /donewith a result /failedwith an error.- Workers loop: claim a job atomically, run it, record the outcome.
The job store¶
# jobstore.py
"""A tiny durable job queue in SQLite. Shared by the API and the worker."""
import json
import sqlite3
import time
import uuid
DB = "jobs.db"
def connect() -> sqlite3.Connection:
conn = sqlite3.connect(DB, timeout=10, isolation_level=None) # autocommit; explicit BEGIN
conn.row_factory = sqlite3.Row
conn.execute("PRAGMA journal_mode=WAL")
return conn
def init() -> None:
with connect() as c:
c.execute("""CREATE TABLE IF NOT EXISTS jobs (
id TEXT PRIMARY KEY,
kind TEXT NOT NULL,
payload TEXT NOT NULL,
idempotency_key TEXT UNIQUE,
status TEXT NOT NULL DEFAULT 'queued', -- queued|running|done|failed
attempts INTEGER NOT NULL DEFAULT 0,
max_attempts INTEGER NOT NULL DEFAULT 3,
run_after REAL NOT NULL DEFAULT 0,
locked_until REAL,
result TEXT,
error TEXT,
created_at REAL NOT NULL)""")
def enqueue(kind: str, payload: dict, idempotency_key: str | None = None) -> tuple[str, bool]:
c = connect()
job_id = uuid.uuid4().hex
try:
c.execute("INSERT INTO jobs (id, kind, payload, idempotency_key, created_at) VALUES (?,?,?,?,?)",
(job_id, kind, json.dumps(payload), idempotency_key, time.time()))
return job_id, True
except sqlite3.IntegrityError:
row = c.execute("SELECT id FROM jobs WHERE idempotency_key = ?", (idempotency_key,)).fetchone()
return row["id"], False
def get(job_id: str) -> dict | None:
row = connect().execute("SELECT * FROM jobs WHERE id = ?", (job_id,)).fetchone()
return dict(row) if row else None
def claim(worker: str, lease_seconds: float = 30) -> dict | None:
"""Atomically take one runnable job. A crashed worker's lease expires and the job is retried."""
now = time.time()
row = connect().execute(
"""UPDATE jobs SET status = 'running', attempts = attempts + 1, locked_until = ?
WHERE id = (SELECT id FROM jobs
WHERE (status = 'queued' AND run_after <= ?)
OR (status = 'running' AND locked_until < ?)
ORDER BY created_at LIMIT 1)
RETURNING *""", (now + lease_seconds, now, now)).fetchone()
return dict(row) if row else None
def finish(job_id: str, result: dict) -> None:
connect().execute("UPDATE jobs SET status='done', result=?, locked_until=NULL WHERE id=?",
(json.dumps(result), job_id))
def fail(job: dict, error: str) -> None:
c = connect()
if job["attempts"] >= job["max_attempts"]:
c.execute("UPDATE jobs SET status='failed', error=?, locked_until=NULL WHERE id=?", (error, job["id"]))
else:
backoff = 2 ** job["attempts"] * 0.25
c.execute("UPDATE jobs SET status='queued', error=?, run_after=?, locked_until=NULL WHERE id=?",
(error, time.time() + backoff, job["id"]))
The important parts:
claimis a singleUPDATE ... WHERE id = (SELECT ... LIMIT 1) RETURNING *. Two workers can't claim the same job, because the database applies the update atomically — no "select, then update" race. (RETURNINGneeds SQLite 3.35 or later; this ran on SQLite 3.53.4. On PostgreSQL the equivalent usesSELECT ... FOR UPDATE SKIP LOCKED.)- Leases (
locked_until): a claimed job belongs to its worker only for a limited time. If the worker dies, the lease expires and another worker picks the job up. - Retries with exponential backoff: a failure re-queues the job with
run_afterin the future, untilmax_attempts. - Idempotency keys: a
UNIQUEcolumn, so a client retrying aPOSTgets the original job back instead of a duplicate. - WAL mode lets the API read while a worker writes.
The API¶
# api.py
import json
from typing import Annotated
from fastapi import FastAPI, Header, HTTPException, Request, Response, status
from pydantic import BaseModel, Field
import jobstore
jobstore.init()
app = FastAPI()
class ReportRequest(BaseModel):
year: int = Field(ge=2000, le=2100)
fail_times: int = Field(0, ge=0, le=5) # for the demo: make the job fail N times first
@app.post("/reports", status_code=status.HTTP_202_ACCEPTED)
def request_report(body: ReportRequest, request: Request, response: Response,
idempotency_key: Annotated[str | None, Header()] = None):
job_id, created = jobstore.enqueue("sales_report", body.model_dump(), idempotency_key)
url = str(request.url_for("job_status", job_id=job_id))
response.headers["Location"] = url
if not created:
response.status_code = status.HTTP_200_OK
return {"job_id": job_id, "status_url": url, "created": created}
@app.get("/jobs/{job_id}")
def job_status(job_id: str, response: Response):
job = jobstore.get(job_id)
if job is None:
raise HTTPException(404, "Job not found")
if job["status"] in ("queued", "running"):
response.headers["Retry-After"] = "1"
return {"id": job["id"], "status": job["status"], "attempts": job["attempts"],
"result": json.loads(job["result"]) if job["result"] else None,
"error": job["error"]}
The worker¶
# worker.py
import json, os, signal, sys, time
import jobstore
running = True
def stop(*_):
global running
running = False
signal.signal(signal.SIGTERM, stop)
NAME = f"worker-{os.getpid()}"
def sales_report(payload: dict, attempt: int) -> dict:
if attempt <= payload.get("fail_times", 0):
raise RuntimeError(f"upstream timeout (simulated, attempt {attempt})")
time.sleep(1.0) # the slow part
return {"year": payload["year"], "total_cents": 123_456_00, "by": NAME}
HANDLERS = {"sales_report": sales_report}
print(f"{NAME} started", flush=True)
while running:
job = jobstore.claim(NAME, lease_seconds=float(os.environ.get("LEASE_SECONDS", "30")))
if job is None:
time.sleep(0.2)
continue
print(f"{NAME} took {job['id'][:8]} attempt {job['attempts']}", flush=True)
try:
result = HANDLERS[job["kind"]](json.loads(job["payload"]), job["attempts"])
jobstore.finish(job["id"], result)
print(f"{NAME} finished {job['id'][:8]}", flush=True)
except Exception as exc:
jobstore.fail(job, str(exc))
print(f"{NAME} failed {job['id'][:8]}: {exc}", flush=True)
print(f"{NAME} stopped cleanly", flush=True)
It runs as its own process — python worker.py — as many copies as you want. It handles
SIGTERM by finishing the current job and exiting, so deploys don't abandon work.
Running it¶
The API under Uvicorn and two workers, all on one machine:
Request a report, then send the same request again with the same idempotency key:
HTTP/1.1 202 Accepted
location: http://localhost:8718/jobs/97e0a08d604b476d9554247bf2871ce9
{"job_id":"97e0a08d604b476d9554247bf2871ce9","status_url":"http://localhost:8718/jobs/97e0a08d604b476d9554247bf2871ce9","created":true}
{"job_id":"97e0a08d604b476d9554247bf2871ce9",...,"created":false} [200]
Same job ID, created: false, status 200 — no duplicate report. A second job was
submitted that fails twice before succeeding (fail_times: 2). Polling both once a
second:
t+0s
retry-after: 1
{"id":"97e0a08d...","status":"running","attempts":1,"result":null,"error":null}
{"id":"8c25bc18...","status":"queued","attempts":1,"result":null,"error":"upstream timeout (simulated, attempt 1)"}
t+1s
{"id":"97e0a08d...","status":"done","attempts":1,"result":{"year":2026,"total_cents":12345600,"by":"worker-74857"},"error":null}
{"id":"8c25bc18...","status":"queued","attempts":2,"result":null,"error":"upstream timeout (simulated, attempt 2)"}
t+2s
{"id":"8c25bc18...","status":"running","attempts":3,...}
t+3s
{"id":"8c25bc18...","status":"done","attempts":3,"result":{"year":2025,"total_cents":12345600,"by":"worker-74859"},"error":"upstream timeout (simulated, attempt 2)"}
(IDs shortened; the first job's later identical lines omitted.) The workers' logs:
worker-74857 started
worker-74857 took 97e0a08d attempt 1
worker-74857 finished 97e0a08d
worker-74859 started
worker-74859 took 8c25bc18 attempt 1
worker-74859 failed 8c25bc18: upstream timeout (simulated, attempt 1)
worker-74859 took 8c25bc18 attempt 2
worker-74859 failed 8c25bc18: upstream timeout (simulated, attempt 2)
worker-74859 took 8c25bc18 attempt 3
worker-74859 finished 8c25bc18
The two jobs ran in parallel on different workers; the failing one was retried after
growing delays and succeeded on its third attempt. The finished job keeps its last error
in the error column — useful for diagnosis; clear it in finish if you'd rather not
show it.
Worked example: a worker dies mid-job¶
A worker was started with a 2-second lease (LEASE_SECONDS=2), given a job, and killed
with kill -9 half a second into the 1-second task:
killed 74903 mid-job
{"id":"4e06b5f8...","status":"running","attempts":1,"result":null,"error":null}
The job was stuck in running — nobody was working on it. A new worker was started;
once the lease expired, it claimed the job:
{"id":"4e06b5f8...","status":"done","attempts":2,"result":{"year":2024,"total_cents":12345600,"by":"worker-74914"},"error":null}
worker-74903 started
worker-74903 took 4e06b5f8 attempt 1
worker-74914 started
worker-74914 took 4e06b5f8 attempt 2
worker-74914 finished 4e06b5f8
worker-74914 stopped cleanly
The last line is the graceful SIGTERM exit. The crash recovery has a consequence you
must design for: the killed worker might have done part of the job — sent some emails,
charged a card — before dying, and the retry does it again. Queues give you
at-least-once execution. Make jobs idempotent (check "already sent?" before
sending; pass an idempotency key to payment providers) so running twice is harmless. Set
the lease longer than the longest normal run, or extend it periodically from the worker
(a "heartbeat"); a lease that expires while a slow job is still running causes a
duplicate run.
How It Actually Works¶
The queue is just rows with a state machine:
queued ──claim──▶ running ──finish──▶ done
▲ │
└──fail (retry)───┤
└──fail (no attempts left)──▶ failed
running with expired lease ──claim──▶ running (attempt n+1)
Correctness rests on the claim being one atomic statement. SQLite serialises writers, so
two workers' UPDATEs run one after the other; the second sees the job already
running with a fresh lease, the subquery picks a different job or none. PostgreSQL
achieves the same with row locks: FOR UPDATE SKIP LOCKED lets many workers claim
different rows concurrently without blocking each other.
Brokered systems implement the same ideas differently: Celery's "acks late" plus a
visibility timeout is a lease; retries with countdowns are run_after; result backends
hold the result column. Knowing the mechanics makes their configuration options make
sense.
Common mistakes¶
- Long work in the request, so proxies time out and clients retry — starting the work twice.
BackgroundTasksfor work that must not be lost.- Select-then-update claiming, letting two workers take the same job.
- No leases: a crashed worker's job stays
runningforever. - Non-idempotent jobs under at-least-once delivery.
- Polling without
Retry-Afteror backoff, hammering the status endpoint. - Unbounded retries — a poison job that always fails, retried forever. Cap attempts
and alert on
failed.
Exercise¶
- Add
max_attemptsand the backoff base to settings, and exposeGET /jobs?status=failedfor operators, with aPOST /jobs/{id}/retrythat re-queues a failed job. - Add a heartbeat: the worker extends
locked_untilevery few seconds while a job runs. Prove with a 10-second job and a 3-second lease that it runs only once. - Push job completion to clients with Server-Sent Events (Level 3 lesson 6) instead of polling.
- Port the store to PostgreSQL with
FOR UPDATE SKIP LOCKED, or replace it with ARQ or Celery, and map each concept above (lease, retry, idempotency, result) to its option in that tool.