Featured Project5 min read

Distributed Task Queue

Celery + RabbitMQ — Circuit Breakers, Observability, Terraform

Production-grade distributed task queue processing asynchronous workloads (image transformation, PDF rendering, webhook dispatch) with FastAPI, Celery, and RabbitMQ. Includes fault-tolerant resilience patterns, full observability, and Terraform-based AWS deployment.

FastAPICeleryRabbitMQRedisPrometheusGrafanaDocker ComposeTerraformAWS ECS FargatePostgreSQL

01 — Problem

What was hard about this

Async workloads — image transforms, PDF rendering, webhook delivery — need queue-based processing. Production task queues fail in subtle ways: silent message drops under broker restarts, duplicate execution when a worker dies mid-task, and cascading failures when a downstream API gets slow and consumes every worker thread. A naive Celery setup hits all three within the first month of real traffic.

02 — Architecture

How the pieces fit

Loading diagram…
The producer checks a Postgres-enforced idempotency key before publishing. RabbitMQ retries with exponential backoff; the worker's failure hook — not the broker — moves terminal failures to a DLQ and records the state transition in Postgres. Redis-backed circuit breakers shed load when downstream APIs slow down.

03 — Decisions

Trade-offs I'd defend in an interview

01Dead-lettering in the application, not the broker

RabbitMQ will dead-letter for you: set x-dead-letter-exchange on a queue and the broker routes rejected messages itself. I chose not to let it. Celery's failure hook classifies the exception, writes a dead_lettered status to Postgres, and forwards a summary to a dedicated queue. That makes the DLQ a SQL predicate rather than a FIFO queue — filterable, paginated, carrying per-attempt audit history, and replayable with a single HTTP call. The cost is real and I'd name it in an interview: broker-level dead-lettering keeps working when my application or my database is down, and mine does not.

02Idempotency keys, not natural-key dedup

Natural-key dedup (e.g. 'has this user_id+image_id been processed?') breaks down when retries cross worker boundaries. Clients pass an idempotency key on enqueue and a uniqueness constraint on the jobs table enforces it, so a duplicate submission returns the original job with a 200 rather than creating a second one. Execution-level dedup is a separate mechanism: the worker skips a job only if it is already completed, and that status is written after the work finishes, never before — marking it first would turn a mid-task crash into silent loss.

03Circuit breakers in Redis, not in-process

An in-process circuit breaker (e.g. pybreaker) doesn't share state across worker processes — each one has to fail independently before tripping. Backing the breaker state in Redis makes it cluster-wide: one worker tripping the breaker protects the whole pool from hammering a sick upstream.

04RabbitMQ over Kafka or SQS

Kafka's strengths (high-throughput log, replay) didn't match the workload — these are jobs, not events. SQS is fine but FIFO queue limits made it awkward. RabbitMQ gives me per-queue priority support (x-max-priority), first-class Celery integration, and a well-understood operational model.

04 — Outcomes

What shipped

  • Every terminal failure landed in an inspectable dead_lettered state — no failure exited the system unrecorded
  • Zero duplicate executions across 10K jobs with workers SIGKILL-ed mid-task
  • At-least-once execution with a narrow residual window, not exactly-once — the duplicate-counting harness cannot certify zero loss
  • 8 services orchestrated via Docker Compose; Terraform deploys the same topology to ECS Fargate

05 — Next

What I'd do if this had another sprint

  • Replace Docker Compose with EKS to practice K8s operational patterns
  • Add OpenTelemetry distributed tracing — correlation IDs are propagated but not yet emitted as spans
  • Add a Kafka topic for fan-out events (e.g. job completed → downstream consumers)
  • Publish load-test results from k6 with p50/p95/p99 graphs as part of the README

06 — Visual proof

See it in code