A high-throughput distributed task queue supporting 5,000+ tasks/sec with exactly-once delivery semantics, consistent hashing, and backpressure-aware dispatching.
Producers (API/CLI) → Dispatcher (consistent hashing + backpressure)
↓
Redis Streams (8 partitions)
↓
Worker Pool (consumer groups)
↓
PostgreSQL (journal + idempotency)
Key properties:
- Exactly-once delivery — Redis Streams consumer groups + PostgreSQL idempotency keys eliminate duplicate execution
- Consistent hashing — 150 virtual nodes per partition for even distribution with minimal remapping on scale events
- Backpressure-aware — monitors pending entry list (PEL) per partition, reroutes or throttles when overloaded
- Automatic failure recovery — heartbeat-based liveness detection reassigns tasks from dead workers within seconds
- p99 latency < 200ms — connection pooling + async I/O + backpressure keep tail latency low under sustained load
# Start everything
docker-compose up --build -d
# Open dashboard
open http://localhost:8000
# Submit a task
curl -X POST http://localhost:8000/api/tasks \
-H "Content-Type: application/json" \
-d '{"task_type": "image.thumbnail", "payload": {"url": "https://picsum.photos/800/600"}}'
# Check task status
curl http://localhost:8000/api/tasks/{task-id}
# View metrics
curl http://localhost:8000/api/metrics
# API docs
open http://localhost:8000/docspip install -r requirements.txt
# Submit task
python -m src.cli.main submit --type image.thumbnail --payload '{"url": "https://picsum.photos/800/600"}'
# Check status
python -m src.cli.main status <task-id>
# List workers
python -m src.cli.main workers
# View metrics
python -m src.cli.main metrics# Run full benchmark suite
make benchmark
# Locust load test only (50 concurrent producers, 60s)
cd benchmarks && python3 -m locust -f locustfile.py --headless -u 50 -r 10 --run-time 60s --host http://localhost:8000
# Verify exactly-once (default 1000 tasks)
python3 benchmarks/verify_exactly_once.py
# Verify exactly-once (10K tasks)
python3 benchmarks/verify_exactly_once.py 10000
# Failure recovery test
python3 benchmarks/failure_recovery.py| Component | Technology | Why |
|---|---|---|
| Task dispatch | Redis Streams | Sub-ms latency, built-in consumer groups |
| Durable journal | PostgreSQL | ACID, idempotency via UNIQUE constraints |
| API server | FastAPI + Uvicorn | Native async, WebSocket, auto-docs |
| Workers | async Python | Event-loop based, non-blocking I/O |
| Dashboard | Vanilla JS + Chart.js | Zero build step, served as static files |
| Image processing | Pillow | Lightweight, covers resize/thumbnail/metadata |
| Load testing | Locust | Python-native, scriptable, good reporting |
| Containers | Docker Compose | Single command deployment |
├── src/
│ ├── config.py # Configuration from env vars
│ ├── dispatcher/ # Hash ring, backpressure, dispatch logic
│ ├── worker/ # Worker loop, heartbeat, liveness, claimer
│ ├── storage/ # PostgreSQL + Redis Streams wrappers
│ ├── api/ # FastAPI server, REST routes, WebSocket
│ ├── cli/ # Click CLI
│ └── handlers/ # Task handlers (image processing)
├── dashboard/ # Static HTML/JS/CSS dashboard
├── benchmarks/ # Locust, verification, recovery tests
├── tests/ # Unit + integration tests
├── docker-compose.yml
├── Dockerfile
└── Makefile