Skip to content

Repository files navigation

Distributed Task Queue System

A production-ready distributed task queue system built with Go, demonstrating advanced microservices patterns including saga orchestration, job dependency management, and real-time monitoring with interactive visualization.

ArchitectureGoPostgreSQLRabbitMQDocker

Features

Core Capabilities

  • Distributed Job Processing: Horizontally scalable worker pool with fair job distribution
  • Job Dependency Chains: Build complex workflows (sequential, parallel, diamond patterns)
  • Saga Pattern: Distributed transaction management with automatic compensation and rollback
  • gRPC + HTTP APIs: High-performance RPC and RESTful endpoints
  • Real-time Dashboard: WebSocket-powered live monitoring with interactive graph visualization
  • Reliable Message Queue: RabbitMQ with dead letter queue and configurable retry logic
  • Persistent Storage: PostgreSQL with optimized indexes for high-throughput operations

Advanced Features

  • Job Chaining: Create complex DAG workflows with dependency tracking
  • Saga Orchestration: Backward recovery pattern with automatic compensation task generation
  • Real-time Visualization: Interactive dependency graphs with live status updates
  • Progress Monitoring: Live logging and percentage-based progress tracking
  • Graceful Degradation: Handles deleted jobs, worker failures, and network partitions
  • Zero-Retry Jobs: Support for fail-fast jobs with max_retries=0
  • Exponential Backoff: Smart retry logic with configurable delays (2^n seconds, capped at 60s)

Developer Experience

  • One-Command Deploy: Full stack up with docker-compose up
  • Interactive Examples: Pre-built workflow patterns (Sequential, Parallel, Diamond, Saga)
  • Clean JSON API: No escaped strings, accepts native JSON objects
  • Comprehensive Docs: Architecture guide, developer docs, and testing scenarios

Architecture

┌─────────────────────────────────────────────────────────────────┐
│ Web Browser │
│ ┌──────────────────────────────────────────────────────────┐ │
│ │ Web Dashboard (HTML/JS) │ │
│ │ • Job submission & chain creation │ │
│ │ • Real-time status table with live updates │ │
│ │ • Interactive dependency graph (vis.js) │ │
│ │ • Live logging panel with progress bars │ │
│ │ • Statistics & metrics dashboard │ │
│ └──────────────────────────────────────────────────────────┘ │
└─────────────────────────────────────────────────────────────────┘
│ │ │
│ HTTP/REST │ WebSocket │ gRPC
▼ ▼ ▼
┌────────────────────────────────────────────────────────────────┐
│ API Service │
│ ┌──────────────┐ ┌──────────────┐ ┌──────────────────────┐ │
│ │ HTTP │ │ WebSocket │ │ gRPC Server │ │
│ │ Handler │ │ Manager │ │ (JobService) │ │
│ │ • Jobs │ │ • Broadcast │ │ • SubmitJob │ │
│ │ • Chains │ │ • Live logs │ │ • SubmitJobChain │ │
│ │ • Sagas │ │ • Metrics │ │ • GetChainStatus │ │
│ └──────────────┘ └──────────────┘ └──────────────────────┘ │
└────────────────────────────┼───────────────────────────────────┘
│
┌───────────┼───────────┐
│ │ │
▼ ▼ ▼
┌──────────┐ ┌──────────┐ ┌──────────┐
│PostgreSQL│ │ RabbitMQ │ │WebSocket │
│ • Jobs │ │ • Main Q │ │ Clients │
│ • Chains │ │ • DLQ │ │ │
│ • Deps │ │ │ │ │
└────┬─────┘ └────┬─────┘ └──────────┘
│ │
│ │ Consume
│ ▼
│ ┌───────────────────┐
│ │ Worker 1 │
│ │ • Task execution │
│ │ • Dependency mgmt │
│ │ • Saga coordinator│
│ └───────────────────┘
│ ║
└────────────╫─────────┐
║ │
┌──────────────┐ │
│ Worker N │ │
└──────────────┘ │
║ │
╚═════════╛
Update Status

Components

  1. API Service (Port 8080 HTTP, 50051 gRPC)

    • Validates jobs, chains, and detects circular dependencies
    • Persists to PostgreSQL with transaction support
    • Publishes jobs to RabbitMQ with message persistence
    • Manages WebSocket connections for real-time updates
    • Serves web dashboard with interactive visualization
  2. Worker Pool (Horizontally Scalable)

    • Consumes jobs from RabbitMQ with fair distribution (prefetch=1)
    • Executes configurable task types with progress reporting
    • Manages job dependencies and triggers child jobs
    • Coordinates saga compensation on failure
    • Handles graceful shutdown (finishes current job on SIGTERM)
  3. PostgreSQL Database

    • Jobs table: metadata, status, results, retry info, compensation data
    • Job chains table: workflow metadata, progress tracking
    • Job dependencies table: parent-child relationships (DAG)
    • Optimized indexes for high-throughput queries
  4. RabbitMQ Message Broker

    • Main job queue with persistent messages
    • Dead letter queue for failed jobs
    • Manual acknowledgment for at-least-once delivery
    • Proper message handling for deleted jobs
  5. Web Dashboard

    • Real-time job table with WebSocket updates
    • Interactive chain visualization (vis.js, left-to-right layout)
    • Live logging panel with progress monitoring
    • Pre-built workflow examples (Sequential, Parallel, Diamond, Saga)

Quick Start

Prerequisites

  • Docker and Docker Compose
  • (Optional) Go 1.24+ for local development

Run with Docker Compose

# Clone the repository
git clone <repository-url>cd distributed-task-queue-system
# Start all services (PostgreSQL, RabbitMQ, API, 3x Workers)
docker-compose up --build
# Access the dashboard
open http://localhost:8080

That's it! The system will start:

  • PostgreSQL on port 5432
  • RabbitMQ on ports 5672 (AMQP) and 15672 (Management UI)
  • API service on ports 50051 (gRPC) and 8080 (HTTP/WebSocket)
  • 3 worker instances (horizontally scalable)

Scale Workers

# Scale to 5 workers
docker-compose up --scale worker=5
# Scale to 10 workers
docker-compose up --scale worker=10

Usage Examples

1. Single Job Submission

Via Web Dashboard:

  1. Navigate to http://localhost:8080
  2. Select job type (sleep, random_fail, word_count, etc.)
  3. Enter payload as JSON object: {"seconds": 5}
  4. Set max_retries (0 = fail immediately, 3 = retry up to 3 times)
  5. Click "Submit Job"
  6. Watch real-time status updates

2. Job Chain (Sequential Pipeline)

Execute jobs in order: A → B → C → D

{
"name": "Data Processing Pipeline",
"stop_on_failure": true,
"jobs": [
{
"job_id": "extract",
"type": "word_count",
"payload": {"text": "raw data from source"},
"max_retries": 3,
"depends_on_job_ids": []
},
{
"job_id": "transform",
"type": "sleep",
"payload": {"seconds": 2},
"max_retries": 3,
"depends_on_job_ids": ["extract"]
},
{
"job_id": "validate",
"type": "word_count",
"payload": {"text": "validate processed data"},
"max_retries": 3,
"depends_on_job_ids": ["transform"]
},
{
"job_id": "load",
"type": "sleep",
"payload": {"seconds": 1},
"max_retries": 3,
"depends_on_job_ids": ["validate"]
}
]
}

Try it:

  • Click "Load Sequential Example" in the dashboard
  • Click "Create Chain"
  • Watch jobs execute in order with graph visualization

3. Job Chain (Parallel Fan-Out)

Root job triggers multiple children simultaneously: Root → [A, B, C]

{
"name": "Parallel Processing",
"stop_on_failure": false,
"jobs": [
{
"job_id": "root",
"type": "word_count",
"payload": {"text": "distribute work"},
"max_retries": 3,
"depends_on_job_ids": []
},
{
"job_id": "child1",
"type": "sleep",
"payload": {"seconds": 2},
"max_retries": 3,
"depends_on_job_ids": ["root"]
},
{
"job_id": "child2",
"type": "sleep",
"payload": {"seconds": 2},
"max_retries": 3,
"depends_on_job_ids": ["root"]
},
{
"job_id": "child3",
"type": "sleep",
"payload": {"seconds": 2},
"max_retries": 3,
"depends_on_job_ids": ["root"]
}
]
}

Try it:

  • Click "Load Parallel Example" in the dashboard
  • Observe all 3 children start simultaneously after root completes

4. Job Chain (Diamond Pattern)

Parallel branches merge: Start → [A, B] → Merge

{
"name": "Diamond Pattern",
"stop_on_failure": true,
"jobs": [
{
"job_id": "start",
"type": "word_count",
"payload": {"text": "start processing"},
"max_retries": 3,
"depends_on_job_ids": []
},
{
"job_id": "branch1",
"type": "sleep",
"payload": {"seconds": 2},
"max_retries": 3,
"depends_on_job_ids": ["start"]
},
{
"job_id": "branch2",
"type": "sleep",
"payload": {"seconds": 3},
"max_retries": 3,
"depends_on_job_ids": ["start"]
},
{
"job_id": "merge",
"type": "word_count",
"payload": {"text": "merge results"},
"max_retries": 3,
"depends_on_job_ids": ["branch1", "branch2"]
}
]
}

Try it:

  • Click "Load Diamond Example" in the dashboard
  • Watch merge job wait for both branches to complete

5. Saga Pattern (Distributed Transaction)

Automatic rollback on failure with compensation tasks.

What is a Saga? A saga is a sequence of local transactions with compensating actions. If any step fails, the system automatically executes compensation tasks in reverse order to undo completed work.

Example: Travel Booking

{
"name": "Book Travel with Compensation",
"is_saga": true,
"saga_policy": "backward_recovery",
"stop_on_failure": true,
"jobs": [
{
"type": "reserve_payment",
"payload": {"amount": 500.00, "customer_id": "cust_123"},
"compensation_type": "release_payment",
"max_retries": 3,
"depends_on_job_ids": []
},
{
"type": "book_flight",
"payload": {"flight": "AA100", "seats": 2},
"compensation_type": "cancel_flight",
"max_retries": 3,
"depends_on_job_ids": ["reserve_payment"]
},
{
"type": "book_hotel",
"payload": {"hotel": "Marriott", "nights": 3},
"compensation_type": "cancel_hotel",
"max_retries": 3,
"depends_on_job_ids": ["book_flight"]
}
]
}

How it works:

  1. Success Path: reserve_payment → book_flight → book_hotel
  2. Failure Path (if book_hotel fails):
    • System detects saga failure
    • Auto-generates compensation jobs in reverse order
    • Executes: cancel_flight → release_payment
    • Compensation tasks receive original job results (e.g., payment amount)
    • System rolls back to initial state

Try it:

  • Click "Load Saga Example" in the dashboard
  • All tasks succeed: Normal flow
  • Modify to trigger failure: Watch automatic compensation

Task Types

Built-in Tasks

sleep - Simulates long-running task

{"type": "sleep", "payload": {"seconds": 5}, "max_retries": 3}

random_fail - Demonstrates retry logic (50% failure rate)

{"type": "random_fail", "payload": {"threshold": 0.5}, "max_retries": 3}

word_count - Counts words and characters

{"type": "word_count", "payload": {"text": "sample text"}, "max_retries": 3}

image_processing - Placeholder for image operations

{"type": "image_processing", "payload": {"image_url": "...", "operations": [...]}, "max_retries": 3}

Saga Forward Tasks

reserve_payment → Compensation: release_payment

{
"type": "reserve_payment",
"payload": {"amount": 500.00, "customer_id": "cust_123"},
"compensation_type": "release_payment"
}

book_flight → Compensation: cancel_flight

{
"type": "book_flight",
"payload": {"flight": "AA100", "seats": 2},
"compensation_type": "cancel_flight"
}

book_hotel → Compensation: cancel_hotel

{
"type": "book_hotel",
"payload": {"hotel": "Marriott", "nights": 3},
"compensation_type": "cancel_hotel"
}

reserve_inventory → Compensation: release_inventory

{
"type": "reserve_inventory",
"payload": {"product_id": "prod_123", "quantity": 10},
"compensation_type": "release_inventory"
}

API Reference

HTTP Endpoints

MethodEndpointDescription
POST/api/jobsSubmit single job
GET/api/jobsList all jobs (with filters)
GET/api/jobs/{id}Get job details
DELETE/api/jobs/allDelete all jobs and purge queue
POST/api/chainsCreate job chain with dependencies
GET/api/chains/{id}Get chain status and all jobs
GET/api/statsGet system statistics
GET/wsWebSocket endpoint for real-time updates

gRPC Methods

Defined in proto/jobs.proto:

  • SubmitJob - Submit single job
  • GetJobStatus - Get job by ID
  • ListJobs - List jobs with filters
  • SubmitJobChain - Submit job chain
  • GetChainStatus - Get chain with all jobs

About

Distributed task queue in Go — RabbitMQ + Postgres, at-least-once delivery, dead-letter queue, saga compensation, and job-dependency graphs (sequential/parallel/diamond) with a live WebSocket dashboard.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages