Skip to content

About

Distributed task queue — FastAPI + Redis broker + PostgreSQL; async workers, priority queues, retry with backoff, dead-letter queue, scheduled jobs, real-time WebSocket dashboard

Resources

Contributing

Stars

0 stars

Watchers

1 watching

Forks

Repository files navigation

Distributed Task Queue

A production-grade distributed task queue system built with FastAPI, Redis, and PostgreSQL. Features async task processing, priority queues, retry logic, dead letter queues, real-time WebSocket monitoring, and a rich dashboard.

Architecture

┌──────────┐     ┌──────────┐     ┌──────────┐
│  Client  │────▶│   API    │────▶│  Redis   │
└──────────┘     │ (FastAPI)│     │  Broker  │
                 └──────────┘     └────┬─────┘
                      │                │
                      ▼                ▼
                 ┌──────────┐     ┌──────────┐
                 │PostgreSQL│     │ Workers  │
                 └──────────┘     └──────────┘
                      │                │
                      ▼                ▼
                 ┌──────────┐     ┌──────────┐
                 │Dashboard │◀────│WebSocket │
                 │ (HTML/JS)│     │  Events  │
                 └──────────┘     └──────────┘

Quick Start

# Clone and start
git clone <repo>
cd distributed_task_queue
docker-compose up -d

# Apply migrations
docker-compose exec app alembic upgrade head

# Seed demo data
docker-compose exec app python scripts/seed_data.py

# Open dashboard
open http://localhost:8080

Manual Setup

python -m venv venv
source venv/bin/activate
pip install -r requirements.txt

# Start dependencies
docker run -d -p 6379:6379 redis:7-alpine
docker run -d -p 5432:5432 -e POSTGRES_PASSWORD=postgres postgres:15-alpine

# Initialize
cp .env.example .env
python scripts/init_db.py
python scripts/seed_data.py

# Start services
uvicorn backend.main:app --reload        # API on :8000
python scripts/start_worker.py           # Worker
python scripts/scheduler_runner.py       # Scheduler

API Endpoints

Method Path Description
POST /api/v1/jobs Create a job
GET /api/v1/jobs List jobs (paginated, filterable)
GET /api/v1/jobs/{id} Get job details
PATCH /api/v1/jobs/{id} Update a job
DELETE /api/v1/jobs/{id} Delete a job
POST /api/v1/jobs/{id}/cancel Cancel a job
GET /api/v1/workers List workers
GET /api/v1/queues List queues and lengths
GET /api/v1/stats System statistics
GET /api/v1/health Health check
WS /api/v1/ws/{channel} WebSocket for real-time updates

Project Structure

├── backend/           # FastAPI application
│   ├── api/           # Routes, schemas, dependencies
│   ├── broker/        # Redis message queue + DLQ
│   ├── models/        # SQLAlchemy ORM models
│   ├── services/      # Business logic layer
│   ├── workers/       # Task executor and pool
│   ├── websocket/     # Real-time event system
│   └── utils/         # Helpers, decorators, serializers
├── config/            # Application configuration
├── dashboard/         # Web-based monitoring UI
├── migrations/        # Alembic database migrations
├── scripts/           # CLI utilities
└── tests/             # Unit, integration, and load tests

Task Handlers

Task Type Description
echo Returns the input payload
sleep Sleeps for a configurable duration
math Performs arithmetic operations
fail Always fails (testing)
http_request Makes HTTP requests

Tech Stack

  • Python 3.11+ / FastAPI - Web framework
  • PostgreSQL + SQLAlchemy async - Persistence
  • Redis - Message broker and queue
  • Alembic - Database migrations
  • APScheduler - Job scheduling
  • WebSocket - Real-time events
  • Docker Compose - Container orchestration

About

Distributed task queue — FastAPI + Redis broker + PostgreSQL; async workers, priority queues, retry with backoff, dead-letter queue, scheduled jobs, real-time WebSocket dashboard

Resources

Contributing

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Used by

Contributors

Languages