Skip to content

Latest commit

 

History

1 Commit

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

Distributed Task Engine

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.

CI Java 21 Spring Boot 3.5 License MIT

▶ Live dashboard demo — the real dashboard running against an in-browser simulated engine (no server needed). For the full system, see Quick start.

What is this

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.

Architecture

                 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)     │
        └───────────────────────────┘         └──────────────────┘

Key technical decisions

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.

Quick start

git clone https://github.com/addictcode/distributed-task-engine.git
cd distributed-task-engine
docker compose up --build

Open http://localhost:8080, hit ⚡ FLOOD DEMO and watch tasks stream through the pipeline live.

API reference

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.

Configuration

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

Running tests

./mvnw test

Unit 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.

Project structure

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

License

MIT

About

Distributed job queue with priority scheduling, retries, DLQ and a real-time dashboard — a mini-Celery in Java 21 + Spring Boot

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages