A distributed job queue with priority scheduling, retries, a dead-letter queue and a real-time monitoring dashboard — a mini-Celery written from scratch in Java.
▶ Live dashboard demo — the real dashboard running against an in-browser simulated engine (no server needed). For the full system, see Quick start.
Distributed Task Engine accepts jobs over a REST API, queues them by priority in Redis, and executes them on a pool of virtual-thread workers that coordinate safely through PostgreSQL — across any number of application instances. Failed jobs are retried with exponential backoff and parked in a dead-letter queue when their retry budget runs out, and every state transition is streamed over STOMP/WebSocket to a zero-dependency monitoring dashboard. It is the core of systems like Celery or Sidekiq, built from first principles to show how task distribution, locking, retries and liveness detection actually work.
POST /api/v1/tasks
│
▼
┌──────────────────┐ task row (source of truth)
│ REST API │ ─────────────────────────────┐
└────────┬─────────┘ ▼
│ push(id, priority) ┌──────────────┐
▼ │ PostgreSQL │
┌───────────────────────────┐ │ │
│ Redis Sorted Set queue │ │ SELECT ... FOR │
│ score = prio·10¹³ − time │ │ UPDATE SKIP │
└────────────┬──────────────┘ │ LOCKED │
│ ZPOPMAX (Lua, atomic) └───────▲──────┘
▼ │ claim / complete
┌───────────────────────────┐ │
│ Worker Pool │ ───────────────────────────┘
│ N virtual threads │
│ heartbeat → Redis TTL │──── failure ──► retry w/ backoff ──► DLQ
└────────────┬──────────────┘
│ every state change
▼
┌───────────────────────────┐ ┌──────────────────┐
│ STOMP over WebSocket │ ──────► │ Dashboard │
│ /topic/tasks │ │ (vanilla JS) │
└───────────────────────────┘ └──────────────────┘
SELECT … FOR UPDATE SKIP LOCKED instead of application-level locking.
The Redis queue distributes task IDs, but the claim itself happens as a row lock in PostgreSQL. If two workers (possibly on different instances) ever receive the same ID, exactly one acquires the row and flips it to RUNNING; the other sees the lock, skips, and moves on — no distributed lock service, no lease timeouts, no split-brain. The database that already stores the truth also arbitrates it, and a crashed worker's lock disappears with its transaction.
Redis Sorted Sets as the priority queue.
A sorted set gives O(log n) insert and an atomic ZPOPMAX, so "take the most urgent task" is one round-trip with no scanning. The score encodes priority × 10¹³ − epochMillis: the priority term dominates, so CRITICAL always beats LOW, while the negated timestamp makes equal-priority tasks FIFO. The pop runs in a Lua script, making read-and-remove a single atomic step between competing workers.
Exponential backoff for retries.
A task that failed because a downstream service is choking will fail again 500 ms later — hammering it just prolongs the outage. Delays of 10s → 40s → 160s give transient failures room to clear while keeping the first retry fast, and the cap (maxRetries, then DLQ) turns persistent failures into an explicit, inspectable queue instead of an infinite loop.
Worker heartbeats with Redis TTL keys.
Every worker refreshes worker:heartbeat:{id} (TTL 30s) each 15s. If an instance dies, its keys silently expire and a reaper re-queues the tasks it was running — work is never lost to a crashed pod.
git clone https://github.com/addictcode/distributed-task-engine.git
cd distributed-task-engine
docker compose up --buildOpen http://localhost:8080, hit ⚡ FLOOD DEMO and watch tasks stream through the pipeline live.
Base path: /api/v1
| Method | Path | Description | Example |
|---|---|---|---|
| POST | /tasks |
Submit a task | curl -X POST localhost:8080/api/v1/tasks -H 'Content-Type: application/json' -d '{"type":"SIMULATE_HEAVY","priority":"HIGH","maxRetries":3,"payload":{"durationSeconds":4}}' |
| POST | /tasks/bulk |
Submit up to 100 tasks | curl -X POST localhost:8080/api/v1/tasks/bulk -H 'Content-Type: application/json' -d '{"tasks":[{"type":"CUSTOM"},{"type":"GENERATE_REPORT"}]}' |
| GET | /tasks |
List tasks (?status=&type=&priority=&page=&size=) |
curl 'localhost:8080/api/v1/tasks?status=DEAD' |
| GET | /tasks/{id} |
Get one task | curl localhost:8080/api/v1/tasks/{id} |
| DELETE | /tasks/{id} |
Cancel a PENDING task | curl -X DELETE localhost:8080/api/v1/tasks/{id} |
| POST | /tasks/{id}/retry |
Re-queue a DEAD task from the DLQ | curl -X POST localhost:8080/api/v1/tasks/{id}/retry |
| POST | /tasks/retry-dead |
Re-queue every DEAD task | curl -X POST localhost:8080/api/v1/tasks/retry-dead |
| GET | /stats |
Throughput, avg execution time, error rate | curl localhost:8080/api/v1/stats |
| GET | /workers |
Active workers and their current tasks | curl localhost:8080/api/v1/workers |
| POST | /demo/flood |
Submit N random demo tasks | curl -X POST 'localhost:8080/api/v1/demo/flood?count=20' |
Task types: SEND_WEBHOOK (HTTP POST to payload.url), GENERATE_REPORT (simulated CPU work), SIMULATE_HEAVY (sleep + configurable failure rate), CUSTOM (echo).
Priorities: LOW=1, NORMAL=5, HIGH=10, CRITICAL=20.
| Property | Default | Description |
|---|---|---|
app.workers.count |
5 |
Worker threads per instance |
app.heartbeat.interval-ms |
15000 |
Heartbeat refresh period |
app.heartbeat.timeout-ms |
30000 |
Silence after which a worker is reaped |
app.scheduler.poll-interval-ms |
500 |
Queue poll sleep / retry sweep cadence |
app.retry.base-delay-seconds |
10 |
Delay before the first retry |
app.retry.multiplier |
4.0 |
Backoff growth factor |
./mvnw testUnit tests cover backoff math, queue scoring and worker transitions; integration tests spin up real PostgreSQL + Redis with Testcontainers and verify submission, execution, the retry → DLQ flow and priority ordering. Docker must be running.
com.taskengine
├── config/ Typed @ConfigurationProperties, Redis & WebSocket wiring
├── domain/ Task, WorkerNode entities; status/priority/type enums
├── repository/ Spring Data JPA repositories (incl. SKIP LOCKED claim)
├── service/
│ ├── TaskService Submission, cancel, claim/complete, DLQ re-queue
│ ├── TaskScheduler Retry sweep + dead-worker reaper
│ ├── WorkerPool Starts/stops virtual-thread workers
│ ├── Worker Pop → claim → execute → record outcome loop
│ ├── RetryService Exponential backoff, DLQ handoff
│ ├── HeartbeatService Redis TTL heartbeats, worker registry
│ └── handlers/ One handler per task type
├── queue/ RedisTaskQueue (sorted set + Lua pop), DeadLetterQueue
├── api/ REST controllers + DTOs + error handling
└── websocket/ STOMP broadcaster for live dashboard updates
MIT