Skip to content

Latest commit

 

History

13 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

Distributed Task Processing Engine

A high-throughput distributed task queue supporting 5,000+ tasks/sec with exactly-once delivery semantics, consistent hashing, and backpressure-aware dispatching.

Architecture

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

Quick Start

# 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/docs

CLI

pip 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

Benchmarks

# 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

Tech Stack

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

Project Structure

├── 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

About

No description, website, or topics provided.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages