A production-ready, distributed job queue system built in Go, inspired by Celery. V2.4 delivers enterprise-grade features including transactional outbox pattern, Kafka streaming, priority-based processing, and multi-machine coordination with automatic failover.
- Multi-worker Processing: Configurable concurrent workers per machine
- Dual Database Support: PostgreSQL for production, SQLite for development
- Lease-based Recovery: Automatic handling of worker failures
- Retry Handling: Smart retry logic with exponential backoff
- Idempotency Support: Duplicate job prevention with custom keys
- Priority-based Processing: 4-level priority system (0-10 scale). Example :
- Critical (10): Fraud detection, urgent alerts
- High (7): User-facing operations, notifications
- Normal (3): Regular business operations
- Background (0): Cleanup, backups, maintenance
- Burst Job Handling: Handle 10k+ concurrent job submissions
- Ultra-High Throughput: ~140 claims/sec on SQLite (benchmarked); PostgreSQL performance significantly higher
- Hybrid Architecture: Kafka + Database with automatic fallback
- Multi-Machine Coordination: Automatic worker discovery across machines
- Heartbeat System: Real-time health monitoring and dead worker cleanup
- Dynamic Fleet Management: Workers join/leave without downtime
- Fleet Visibility: Real-time monitoring via
/workersAPI - Graceful Shutdown: Clean process termination with proper cleanup
- Prometheus Metrics: Built-in observability and monitoring
- HTTP Producer API: RESTful job submission and management
- Health Checks:
/healthand/readyendpoints - Job Tracking: Complete job lifecycle visibility
- Performance Logging: Per-job timing and fleet statistics
- Kafka Integration: High-throughput message streaming for job distribution
- Transactional Outbox: Reliable message delivery with consistency guarantees
- Zero-Downtime Resilience: Automatic fallback to database polling
graph TB
subgraph "Vesper V2.4 - Hybrid Architecture with Transactional Outbox"
subgraph "Machine A"
A1[HTTP API :8080] --> TX1{Transaction}
TX1 --> DB[(Database)]
TX1 --> OB1[Outbox Events]
OB1 --> R1[Outbox Relay]
R1 --> K[Kafka Topic]
K --> KW1[Kafka Workers 1-3]
DB --> DW1[DB Workers 4-5]
KW1 --> DB
DW1 --> DB
end
subgraph "Machine B"
A2[HTTP API :8081] --> TX2{Transaction}
TX2 --> DB
TX2 --> OB2[Outbox Events]
OB2 --> R2[Outbox Relay]
R2 --> K
K --> KW2[Kafka Workers 6-8]
DB --> DW2[DB Workers 9-10]
KW2 --> DB
DW2 --> DB
end
subgraph "Machine C"
A3[HTTP API :8082] --> TX3{Transaction}
TX3 --> DB
TX3 --> OB3[Outbox Events]
OB3 --> R3[Outbox Relay]
R3 --> K
K --> KW3[Kafka Workers 11-13]
DB --> DW3[DB Workers 14-15]
KW3 --> DB
DW3 --> DB
end
WR[Worker Registry] --> DB
end
H[User/API] --> A1
H --> A2
H --> A3
style A1 fill:#e1f5fe
style A2 fill:#e1f5fe
style A3 fill:#e1f5fe
style TX1 fill:#fff3e0
style TX2 fill:#fff3e0
style TX3 fill:#fff3e0
style OB1 fill:#f3e5f5
style OB2 fill:#f3e5f5
style OB3 fill:#f3e5f5
style R1 fill:#e8f5e8
style R2 fill:#e8f5e8
style R3 fill:#e8f5e8
style KW1 fill:#4caf50
style KW2 fill:#4caf50
style KW3 fill:#4caf50
style DW1 fill:#2196f3
style DW2 fill:#2196f3
style DW3 fill:#2196f3
style K fill:#ff9800
style DB fill:#fff3e0
style WR fill:#f3e5f5
Transactional Outbox Pattern:
- HTTP API writes jobs atomically to both Jobs and Outbox Events tables
- Ensures data consistency and prevents message loss during failures
- Outbox Relay processes events asynchronously with retry logic
Hybrid Worker Distribution:
- Worker Registry: Distributed coordination and heartbeat monitoring
Multi-Machine Coordination:
- Shared database for job state management
- Kafka consumer groups for load balancing
- Automatic worker discovery and failover
- Outbox relay instances on each machine for high availability
sequenceDiagram
participant P as Producer
participant W as Worker
participant DB as Database
Note over P,DB: Job Submission Flow
P->>DB: 1. Create job (status: pending)
Note over DB: Job Processing Flow
W->>DB: 2. Atomically claim job (pending -> processing)
W->>W: 3. Execute job logic
alt Job Success
W->>DB: 4a. Update status (done)
else Job Failure
W->>DB: 4b. Update status (failed)
end
stateDiagram-v2
[*] --> pending: Job Created
pending --> processing: Worker Claims Job
processing --> done: Job Completed Successfully
processing --> failed: Job Failed/Error
done --> [*]
failed --> [*]
note right of pending: Job waiting in queue
note right of processing: Worker executing job
note right of done: Job completed successfully
note right of failed: Job failed with error
- Go 1.19 or higher
- Database: PostgreSQL (production) or SQLite (development)
go get github.com/akshit-git24/vesper-hivemind/vesper/v2@v2.4.0The repo root includes a minimal main.go that registers a consumer and calls Vesper.Run().
For local development, create a workspace file so the root app uses the local ./vesper module instead of the published release:
cp go.work.example go.workgo.work is intentionally gitignored, so GitHub CI and production builds stay independent and resolve the released vesper/v2 version from go.mod.
-
Clone the repository
git clone <repository-url> cd vesper-hivemind
-
Install dependencies
go mod tidy
-
Configure environment
cp .env.example .env # Edit .env with your settings -
Run the app
go run .
Register a consumer function in your own main.go and the library will call it for every job:
package main
import (
vesper "github.com/akshit-git24/vesper-hivemind/vesper/v2"
"github.com/joho/godotenv"
)
func main() {
godotenv.Load()
vesper.RegisterConsumer(func(job vesper.Job) error {
// your business logic here
return nil
})
vesper.Run()
}# Database Configuration
DB_HOST="localhost" # PostgreSQL host (if provided, uses PostgreSQL)
DB_PORT="5432" # PostgreSQL port
DB_USER="vesper" # PostgreSQL username
DB_PASSWORD="secure123" # PostgreSQL password
DB_NAME="vesper" # PostgreSQL database name
DB_SSLMODE="disable" # PostgreSQL SSL mode
DB_TIMEZONE="UTC" # PostgreSQL timezone
# SQLite Fallback (used when DB_HOST is not provided)
DB_NAME="sqlite.db" # SQLite database file path
# Application Configuration
WORKERS="5" # number of concurrent workers
API_ADDR=":8080" # http server address
APP_ENV="development" # or "production"
REQUIRE_CONSUMER="0" # set "1" to require a registered consumer at startup
# Kafka Configuration
KAFKA_ENABLED="true" # enable/disable Kafka integration
KAFKA_BROKERS="localhost:9092" # Kafka broker addresses (comma-separated)
KAFKA_TOPIC="Jobs" # Kafka topic name for job messages
KAFKA_GROUP_ID="vesper-workers" # Kafka consumer group ID
KAFKA_TOPIC_PARTITIONS="3" # number of topic partitions
KAFKA_TOPIC_REPLICATION_FACTOR="1" # topic replication factor
# HTTP Timeouts
HTTP_READ_HEADER_TIMEOUT="5s" # http read header timeout
HTTP_READ_TIMEOUT="10s" # http read timeout
HTTP_WRITE_TIMEOUT="10s" # http write timeout
HTTP_IDLE_TIMEOUT="60s" # http idle timeoutEnvironment behavior:
| APP_ENV | REQUIRE_CONSUMER | Result |
|---|---|---|
| production | anything | Require consumer |
| development | 1 | Require consumer |
| development | 0 / empty | Allow startup |
| Configuration | Database Used | Use Case |
|---|---|---|
DB_HOST provided |
PostgreSQL | Production, multi-machine deployment |
DB_HOST not set |
SQLite | Development, single-machine testing |
- Workers:
5 - PostgreSQL: Auto-detected when
DB_HOSTis provided - SQLite fallback:
Vesper_sqlite.db
Vesper includes production-ready load testing scripts to validate burst performance:
# Standard burst test (100 jobs)
.\test-burst.ps1
# High-load test (1000 jobs, 25 concurrent)
.\test-burst.ps1 -JobCount 1000 -ThrottleLimit 25
# Extreme load test (10,000 jobs)
.\test-10k-burst.ps1
# Custom test with detailed output
.\test-10k-burst.ps1 -JobCount 5000 -DetailedOutputBenchmarks run on Windows, Intel i7-13620H:
| Workers | SQLite (dev) | PostgreSQL (local Docker) |
|---|---|---|
| 5 workers — ClaimNextJob | ~140 claims/sec | ~45–187 claims/sec |
| 5 workers — EndToEnd | ~94–496 jobs/sec | ~225 jobs/sec |
| 30 workers — EndToEnd | ~89 jobs/sec | ~132 jobs/sec |
SQLite degrades under high concurrency (single connection by design). PostgreSQL with SELECT FOR UPDATE SKIP LOCKED scales with worker count and completed 30-worker benchmarks 6x faster than SQLite.
Production PostgreSQL on dedicated hardware will significantly outperform these local Docker numbers.
Run benchmarks yourself:
go test -bench="Benchmark" -benchtime=10s -count=3 -run=^$ .
- Burst Capacity: Handle 10k+ concurrent submissions
- Latency: Sub-100ms job submission response times
- Scalability: Linear scaling with additional machines
- Realistic Job Distribution: Mixed priority levels and job types
- Concurrent Submission: Parallel HTTP requests with throttling
- Performance Metrics: Throughput, success rates, timing analysis
- Error Handling: Retry logic with exponential backoff
- Comprehensive Reporting: Detailed statistics and performance ratings
# Terminal 1
WORKERS=10 API_ADDR=:8080 go run .
# Terminal 2
WORKERS=15 API_ADDR=:8081 go run .
# Terminal 3
WORKERS=5 API_ADDR=:8082 go run .
# Check fleet status
curl http://localhost:8080/workersAutomatic PostgreSQL Detection: When DB_HOST is provided, Vesper automatically uses PostgreSQL for true distributed deployment.
# Setup shared PostgreSQL database
docker run -d --name postgres \
-e POSTGRES_DB=vesper \
-e POSTGRES_USER=vesper \
-e POSTGRES_PASSWORD=secure123 \
-p 5432:5432 postgres:15
# Machine A (192.168.1.100)
DB_HOST=192.168.1.100 DB_USER=vesper DB_PASSWORD=secure123 WORKERS=20 go run .
# Machine B (192.168.1.101)
DB_HOST=192.168.1.100 DB_USER=vesper DB_PASSWORD=secure123 WORKERS=15 go run .
# Machine C (192.168.1.102)
DB_HOST=192.168.1.100 DB_USER=vesper DB_PASSWORD=secure123 WORKERS=10 go run .
# Monitor distributed fleet
curl http://192.168.1.100:8080/workersResult: 45 total workers across 3 machines processing jobs from shared PostgreSQL database with automatic load balancing and failover.
# View all workers
curl http://localhost:8080/workers | jq '.'
# Check total capacity
curl http://localhost:8080/workers | jq '.total_workers'
# Monitor by machine
curl http://localhost:8080/workers | jq '.machines'
# Graceful shutdown (Ctrl+C or SIGTERM)
kill -TERM <process_id>{
"job_id": "550e8400-e29b-41d4-a716-446655440000",
"type": "resize_image",
"priority": 7,
"payload": {
"image_url": "https://example.com/cat.png",
"width": 640,
"height": 480
},
"status": "pending",
"retries": 0,
"error": "",
"created_at": "2024-01-15T10:30:00Z",
"updated_at": "2024-01-15T10:30:00Z"
}0: Background tasks (default) - cleanup, backups3: Regular tasks - normal business operations7: Important tasks - user-facing operations10: Critical tasks - urgent alerts, fraud detection
Note: Any integer 0-10 is valid. Higher numbers = higher priority.
pending: Job created and waiting for processingprocessing: Job currently being executed by a workerdone: Job completed successfullyfailed: Job failed during executiondlq: Job moved to dead-letter queue after max retries
Vesper V2.4 features a production-ready Kafka integration with transactional outbox pattern for guaranteed message delivery and system consistency.
- Ultra-High Performance: ~140 claims/sec on SQLite (benchmarked); PostgreSQL performance significantly higher
- Transactional Outbox: Guarantees message delivery consistency with database transactions
- Automatic Fallback: Seamlessly switches to database-only mode if Kafka is unavailable
- Horizontal Scaling: Add Kafka partitions and workers for linear performance scaling
- Real-time Processing: Workers receive job notifications instantly via message streaming
- Zero Configuration: Topics and consumer groups created automatically
- Reliability: Outbox relay ensures no message loss even during system failures
- Message Loss Prevention: Transactional outbox ensures no jobs are lost during failures
- System Consistency: Database and Kafka stay synchronized through atomic transactions
- Latency Issues: Eliminates polling delays with real-time message delivery
- Scalability Limits: Enables horizontal scaling beyond single database constraints
- High Traffic Bursts: Handles sudden job spikes through Kafka's buffering capabilities
- System Reliability: Maintains operation even during Kafka outages via database fallback
- Job Submission: HTTP API writes job + outbox event in single database transaction
- Outbox Relay: Background process publishes events from outbox to Kafka
- Kafka Distribution: Messages distributed across partitions to worker consumer groups
- Hybrid Architecture: Kafka workers process real-time. Without kafka, DB workers handle fallback
- Status Updates: Job status updated in database regardless of processing path
# Enable Kafka integration
KAFKA_ENABLED="true"
KAFKA_BROKERS="localhost:9092"
KAFKA_TOPIC="Jobs"
KAFKA_GROUP_ID="vesper-workers"
KAFKA_TOPIC_PARTITIONS="3"
KAFKA_TOPIC_REPLICATION_FACTOR="1"For WORKERS=100: Automatically creates 75 Kafka workers + 25 database workers for optimal performance and reliability.
- Accepts
POST /jobs - Creates job records in database
- Multiple concurrent workers (configurable)
- Pulls jobs from DB with atomic claim + lease
- Updates job status during processing
- Executes job logic with error handling
- PostgreSQL: Production-grade database with optimized indexes and JSONB support
- SQLite: Development fallback with automatic detection
- GORM ORM: Database operations with automatic schema migration
- Job status tracking: Complete lifecycle management
- Worker registration: Distributed heartbeat and coordination system
# Background task (priority 0 - default)
curl -X POST http://localhost:8080/jobs \
-H "Content-Type: application/json" \
-d '{
"type": "cleanup_temp_files",
"priority": 0,
"payload": {"directory": "/tmp"}
}'
# Regular task (priority 3)
curl -X POST http://localhost:8080/jobs \
-H "Content-Type: application/json" \
-d '{
"type": "resize_image",
"priority": 3,
"payload": {"image_url": "cat.png", "width": 640}
}'
# Important task (priority 7)
curl -X POST http://localhost:8080/jobs \
-H "Content-Type: application/json" \
-d '{
"type": "send_welcome_email",
"priority": 7,
"payload": {"user_id": 123, "email": "user@example.com"}
}'
# Critical task (priority 10)
curl -X POST http://localhost:8080/jobs \
-H "Content-Type: application/json" \
-d '{
"type": "fraud_alert",
"priority": 10,
"payload": {"transaction_id": "tx_456", "risk_score": 0.95}
}'Processing Order: Critical tasks (10) → Important tasks (7) → Regular tasks (3) → Background tasks (0)
Within the same priority level, jobs are processed FIFO (first-in, first-out).
POST /jobs
Header: Idempotency-Key: <your-key>
{
"type": "resize",
"priority": 7,
"payload": {
"image_url": "https://example.com/cat.png",
"width": 640,
"height": 480
}
}Response: 201 Created with the created job JSON.
GET /jobs/{id}
Returns 200 OK with the job JSON or 404 Not Found.
GET /workers
Returns JSON with all active workers across machines and fleet statistics.
GET /health
Returns 200 OK.
GET /ready
Returns 200 OK if DB is reachable.
GET /metrics
Returns Prometheus text format for scraping.
GET /stats
Returns JSON counters and queue stats.
Vesper exposes Prometheus metrics at GET /metrics. The repo ships with a ready-to-use prometheus.yml for a local Vesper process on :8080.
Normal usage:
-
Start the app with
API_ADDR=":8080". -
Keep the repo root
prometheus.ymlas:global: scrape_interval: 15s evaluation_interval: 15s scrape_configs: - job_name: vesper metrics_path: /metrics static_configs: - targets: - localhost:8080
-
Run Prometheus from the repo root:
prometheus --config.file=prometheus.yml
-
Open
http://localhost:9090/targetsand confirm thevespertarget isUP.
Docker container usage for Prometheus:
-
Keep the Vesper app reachable on the host at
:8080. -
Change the scrape target in
prometheus.ymlto:global: scrape_interval: 15s evaluation_interval: 15s scrape_configs: - job_name: vesper metrics_path: /metrics static_configs: - targets: - host.docker.internal:8080
-
Start Prometheus in Docker and mount the config file:
docker run --rm -p 9090:9090 -v "${PWD}/prometheus.yml:/etc/prometheus/prometheus.yml" prom/prometheus -
Open
http://localhost:9090/targetsand confirm thevespertarget isUP.
localhost:8080 works when Prometheus runs directly on your machine. host.docker.internal:8080 is the safer choice when Prometheus runs in a container on Windows or macOS.
Run tests locally:
cd vesper
go test ./...CI:
- GitHub Actions runs
go test ./...andgo vet ./...insidevesper. - Workflow file:
.github/workflows/ci.yml
- Priority-based Job Processing: 4-level priority system with FIFO within levels
- Burst Job Handling: Handle 10k+ concurrent job submissions
- Distributed Architecture: Multi-machine worker coordination
- Fleet Management: Dynamic worker registration and heartbeat monitoring
- Load Testing Suite: Production-ready performance validation scripts
- Comprehensive API: Job submission, retrieval, worker monitoring
- Production Monitoring: Prometheus metrics and health checks
- Kafka Integration: High-throughput message streaming with hybrid architecture
- Transactional Outbox: Reliable message delivery with consistency guarantees
- Outbox Relay System: Background process for reliable event publishing
- Transactional Outbox Pattern: Guarantees message delivery consistency
- Enhanced Kafka Integration: Production-ready with automatic topic creation
- Outbox Relay System: Background process with retry logic for reliable publishing
- Enhanced Monitoring: Outbox events tracking and cleanup processes
- Zero-Configuration Kafka: Automatic topic and consumer group management
This is a learning project focused on understanding distributed systems concepts. Feel free to:
- Report issues or bugs
- Suggest improvements for V2
- Submit pull requests for bug fixes
- Share feedback on architecture decisions
Apache-2.0 License - see LICENSE file for details.