Skip to content

Latest commit

 

History

History
676 lines (518 loc) · 15.5 KB

File metadata and controls

676 lines (518 loc) · 15.5 KB

Workers Guide

Complete guide to building and running workers with Spooled Cloud.

Table of Contents

  1. Overview
  2. Worker Lifecycle
  3. Building Workers
  4. Worker Registration
  5. Processing Jobs
  6. Concurrency
  7. Health & Heartbeats
  8. Graceful Shutdown
  9. gRPC Workers
  10. Best Practices

Overview

Workers are your applications that process jobs from Spooled queues. Spooled itself doesn't execute your business logic - it manages the queue, and your workers do the actual work.

Key concepts:

  • Workers claim jobs via REST or gRPC API
  • Jobs have leases (default 5 minutes) - complete before expiry
  • Workers send heartbeats to stay registered
  • Multiple workers can process the same queue concurrently

Worker Lifecycle

┌──────────┐     ┌──────────┐     ┌──────────┐
│ Register │────▶│ Healthy  │────▶│ Degraded │
└──────────┘     └────┬─────┘     └────┬─────┘
                      │                │
                      │ heartbeat      │ missed heartbeats
                      │                │
                      ▼                ▼
                 ┌──────────┐    ┌──────────┐
                 │ Healthy  │    │ Offline  │
                 └──────────┘    └──────────┘

Worker States

State Description
healthy Sending heartbeats, ready for jobs
degraded Missed 2+ heartbeats, may be struggling
offline Missed 5+ heartbeats, jobs will be recovered

Building Workers

Simple Polling Worker (REST)

import requests
import time
import os

API_KEY = os.environ["SPOOLED_API_KEY"]
BASE_URL = "https://api.spooled.cloud"
WORKER_ID = f"worker-{os.getpid()}"
QUEUE = "emails"

headers = {
    "Authorization": f"Bearer {API_KEY}",
    "Content-Type": "application/json"
}

def main():
    print(f"Worker {WORKER_ID} starting...")
    
    while True:
        try:
            # Claim jobs
            response = requests.post(
                f"{BASE_URL}/api/v1/jobs/claim",
                headers=headers,
                json={
                    "queue_name": QUEUE,
                    "worker_id": WORKER_ID,
                    "limit": 10
                }
            )
            jobs = response.json().get("jobs", [])
            
            if not jobs:
                time.sleep(1)  # No jobs, wait before polling
                continue
            
            for job in jobs:
                process_job(job)
                
        except Exception as e:
            print(f"Error: {e}")
            time.sleep(5)

def process_job(job):
    job_id = job["id"]
    print(f"Processing {job_id}")
    
    try:
        # Your business logic here
        result = do_work(job["payload"])
        
        # Mark complete
        requests.post(
            f"{BASE_URL}/api/v1/jobs/{job_id}/complete",
            headers=headers,
            json={"worker_id": WORKER_ID, "lease_id": job.get("lease_id"), "result": result}
        )
        print(f"✓ Completed {job_id}")
        
    except Exception as e:
        # Mark failed (will retry)
        requests.post(
            f"{BASE_URL}/api/v1/jobs/{job_id}/fail",
            headers=headers,
            json={"worker_id": WORKER_ID, "lease_id": job.get("lease_id"), "error": str(e)}
        )
        print(f"✗ Failed {job_id}: {e}")

def do_work(payload):
    # Your actual processing logic
    return {"success": True}

if __name__ == "__main__":
    main()

Node.js Worker

const axios = require('axios');

const API_KEY = process.env.SPOOLED_API_KEY;
const BASE_URL = 'https://api.spooled.cloud';
const WORKER_ID = `worker-${process.pid}`;
const QUEUE = 'emails';

const client = axios.create({
  baseURL: BASE_URL,
  headers: { 'Authorization': `Bearer ${API_KEY}` }
});

async function main() {
  console.log(`Worker ${WORKER_ID} starting...`);
  
  while (true) {
    try {
      const { data } = await client.post('/api/v1/jobs/claim', {
        queue_name: QUEUE,
        worker_id: WORKER_ID,
        limit: 10
      });
      
      const jobs = data.jobs || [];
      
      if (jobs.length === 0) {
        await sleep(1000);
        continue;
      }
      
      for (const job of jobs) {
        await processJob(job);
      }
      
    } catch (err) {
      console.error('Error:', err.message);
      await sleep(5000);
    }
  }
}

async function processJob(job) {
  console.log(`Processing ${job.id}`);
  
  try {
    const result = await doWork(job.payload);
    
    await client.post(`/api/v1/jobs/${job.id}/complete`, { result });
    console.log(`✓ Completed ${job.id}`);
    
  } catch (err) {
    await client.post(`/api/v1/jobs/${job.id}/fail`, { reason: err.message });
    console.log(`✗ Failed ${job.id}: ${err.message}`);
  }
}

async function doWork(payload) {
  // Your actual processing logic
  return { success: true };
}

const sleep = (ms) => new Promise(r => setTimeout(r, ms));

main();

PHP Worker

<?php

use Spooled\SpooledClient;
use Spooled\Config\ClientOptions;
use Spooled\Worker\SpooledWorker;
use Spooled\Worker\WorkerConfig;
use Spooled\Worker\JobContext;

$client = new SpooledClient(new ClientOptions(
    apiKey: getenv('SPOOLED_API_KEY'),
));

$worker = new SpooledWorker($client, new WorkerConfig(
    queueName: 'emails',
    concurrency: 10,
    pollIntervalMs: 1000,
));

$worker->process(function (JobContext $ctx): array {
    echo "Processing {$ctx->job->id}\n";
    
    $result = doWork($ctx->payload);
    
    echo "✓ Completed {$ctx->job->id}\n";
    return $result;
});

$worker->on('error', function (\Throwable $e, ?JobContext $ctx): void {
    if ($ctx) {
        echo "✗ Failed {$ctx->job->id}: {$e->getMessage()}\n";
    } else {
        echo "Worker error: {$e->getMessage()}\n";
    }
});

// Graceful shutdown
pcntl_signal(SIGTERM, fn() => $worker->stop());
pcntl_signal(SIGINT, fn() => $worker->stop());

$worker->start();

Worker Registration

For advanced use cases, register your worker to enable monitoring and job routing.

Register Worker

curl -X POST https://api.spooled.cloud/api/v1/workers/register \
  -H "Authorization: Bearer sp_live_YOUR_API_KEY" \
  -H "Content-Type: application/json" \
  -d '{
    "id": "worker-abc123",
    "queue_names": ["emails", "notifications"],
    "max_concurrent_jobs": 10,
    "hostname": "worker-1.example.com",
    "metadata": {
      "version": "1.0.0",
      "region": "us-east-1"
    }
  }'

Response:

{
  "id": "worker-abc123",
  "status": "healthy",
  "registered_at": "2024-01-15T10:00:00Z",
  "heartbeat_interval_seconds": 10
}

Deregister Worker

curl -X POST https://api.spooled.cloud/api/v1/workers/worker-abc123/deregister \
  -H "Authorization: Bearer sp_live_YOUR_API_KEY"

Processing Jobs

Claim Jobs

curl -X POST https://api.spooled.cloud/api/v1/jobs/claim \
  -H "Authorization: Bearer sp_live_YOUR_API_KEY" \
  -H "Content-Type: application/json" \
  -d '{
    "queue_name": "emails",
    "worker_id": "worker-abc123",
    "limit": 10
  }'

Complete Job

curl -X POST https://api.spooled.cloud/api/v1/jobs/job_xxx/complete \
  -H "Authorization: Bearer sp_live_YOUR_API_KEY" \
  -H "Content-Type: application/json" \
  -d '{
    "result": {"message_id": "msg_123", "delivered": true}
  }'

Fail Job

curl -X POST https://api.spooled.cloud/api/v1/jobs/job_xxx/fail \
  -H "Authorization: Bearer sp_live_YOUR_API_KEY" \
  -H "Content-Type: application/json" \
  -d '{
    "worker_id": "worker-1",
    "lease_id": "lease_from_claim",
    "error": "SMTP server unavailable"
  }'

Renew Lease

For long-running jobs, renew the lease before it expires:

curl -X POST https://api.spooled.cloud/api/v1/jobs/job_xxx/renew-lease \
  -H "Authorization: Bearer sp_live_YOUR_API_KEY" \
  -H "Content-Type: application/json" \
  -d '{
    "extension_seconds": 300
  }'

Concurrency

Configuring Concurrency

Control how many jobs a worker processes simultaneously:

import asyncio
from concurrent.futures import ThreadPoolExecutor

CONCURRENCY = 10
executor = ThreadPoolExecutor(max_workers=CONCURRENCY)

async def worker_loop():
    while True:
        jobs = await claim_jobs(limit=CONCURRENCY)
        
        if jobs:
            # Process concurrently
            futures = [
                executor.submit(process_job, job) 
                for job in jobs
            ]
            # Wait for all to complete
            for future in futures:
                future.result()
        else:
            await asyncio.sleep(1)

Concurrency Best Practices

Resource Type Suggested Concurrency
CPU-bound 1-2 per CPU core
I/O-bound (API calls) 10-50
Memory-heavy Based on available RAM
External rate limits Match external limit

Health & Heartbeats

Send Heartbeat

Registered workers should send heartbeats every 10 seconds:

curl -X POST https://api.spooled.cloud/api/v1/workers/worker-abc123/heartbeat \
  -H "Authorization: Bearer sp_live_YOUR_API_KEY" \
  -H "Content-Type: application/json" \
  -d '{
    "status": "healthy",
    "current_jobs": 5,
    "cpu_percent": 45,
    "memory_percent": 60
  }'

Heartbeat Response

{
  "acknowledged": true,
  "next_heartbeat_seconds": 10,
  "commands": []
}

Worker Health Check

curl https://api.spooled.cloud/api/v1/workers/worker-abc123 \
  -H "Authorization: Bearer sp_live_YOUR_API_KEY"

Response:

{
  "id": "worker-abc123",
  "status": "healthy",
  "queue_names": ["emails"],
  "current_jobs": 5,
  "last_heartbeat": "2024-01-15T10:00:00Z",
  "registered_at": "2024-01-15T09:00:00Z"
}

Graceful Shutdown

Handle shutdown signals properly to avoid job loss:

import signal
import sys

shutdown_requested = False

def signal_handler(signum, frame):
    global shutdown_requested
    print("Shutdown signal received, finishing current jobs...")
    shutdown_requested = True

signal.signal(signal.SIGTERM, signal_handler)
signal.signal(signal.SIGINT, signal_handler)

def main():
    while not shutdown_requested:
        jobs = claim_jobs(limit=1)  # Claim fewer jobs when approaching shutdown
        
        for job in jobs:
            if shutdown_requested:
                # Don't start new jobs, let this one finish
                break
            process_job(job)
    
    # Deregister worker
    deregister_worker()
    print("Worker shutdown complete")
    sys.exit(0)

Kubernetes Graceful Shutdown

spec:
  containers:
    - name: worker
      lifecycle:
        preStop:
          exec:
            command: ["/bin/sh", "-c", "sleep 30"]
  terminationGracePeriodSeconds: 60

gRPC Workers

For high-throughput scenarios, use gRPC streaming:

Stream Jobs (Server Streaming)

Jobs are pushed to your worker as they become available:

import grpc
from spooled_pb2 import StreamJobsRequest
from spooled_pb2_grpc import QueueServiceStub

# Spooled Cloud (TLS over 443)
channel = grpc.secure_channel('grpc.spooled.cloud:443', grpc.ssl_channel_credentials())

# Self-hosted / local dev (no TLS by default)
# channel = grpc.insecure_channel('localhost:50051')
stub = QueueServiceStub(channel)

request = StreamJobsRequest(
    queue_name="emails",
    worker_id="worker-1",
    lease_duration_secs=300
)

metadata = [('x-api-key', 'sp_live_YOUR_API_KEY')]

for job in stub.StreamJobs(request, metadata=metadata):
    print(f"Received job: {job.id}")
    process_job(job)
    stub.Complete(CompleteRequest(job_id=job.id, worker_id="worker-1"), metadata=metadata)

Bidirectional Streaming

For maximum control, use bidirectional streaming:

def process_jobs():
    while True:
        # Request jobs
        yield ProcessRequest(
            dequeue=DequeueRequest(
                queue_name="emails",
                worker_id="worker-1",
                batch_size=10
            )
        )
        
        # Receive and process
        # ...
        
        # Send completion
        yield ProcessRequest(
            complete=CompleteRequest(
                job_id=job.id,
                worker_id="worker-1"
            )
        )

for response in stub.ProcessJobs(process_jobs(), metadata=metadata):
    handle_response(response)

Best Practices

1. Use Worker IDs

Always include a unique worker ID for debugging:

WORKER_ID = f"{hostname}-{pid}-{uuid.uuid4().hex[:8]}"

2. Handle Lease Expiry

Monitor job processing time and renew leases:

import threading

def process_with_lease_renewal(job):
    # Start lease renewal thread
    stop_renewal = threading.Event()
    renewal_thread = threading.Thread(
        target=renew_lease_periodically,
        args=(job["id"], stop_renewal)
    )
    renewal_thread.start()
    
    try:
        result = do_work(job["payload"])
        complete_job(job["id"], result)
    finally:
        stop_renewal.set()
        renewal_thread.join()

3. Implement Backpressure

Don't claim more jobs than you can handle:

current_jobs = 0
MAX_CONCURRENT = 10

while True:
    available_slots = MAX_CONCURRENT - current_jobs
    if available_slots > 0:
        jobs = claim_jobs(limit=available_slots)
        current_jobs += len(jobs)
        for job in jobs:
            process_async(job, on_complete=lambda: current_jobs -= 1)
    else:
        time.sleep(0.1)  # Wait for slots to free up

4. Log Everything

import structlog

log = structlog.get_logger()

def process_job(job):
    log.info("processing_job", job_id=job["id"], queue=job["queue_name"])
    
    try:
        result = do_work(job["payload"])
        log.info("job_completed", job_id=job["id"], result=result)
    except Exception as e:
        log.error("job_failed", job_id=job["id"], error=str(e))
        raise

5. Use Health Checks

Implement a /health endpoint for your worker:

from flask import Flask
app = Flask(__name__)

@app.route('/health')
def health():
    return {
        "status": "healthy",
        "worker_id": WORKER_ID,
        "current_jobs": len(active_jobs),
        "uptime_seconds": time.time() - start_time
    }

6. Monitor Worker Metrics

Expose Prometheus metrics:

from prometheus_client import Counter, Gauge, Histogram

jobs_processed = Counter('worker_jobs_processed_total', 'Jobs processed', ['queue', 'status'])
jobs_in_progress = Gauge('worker_jobs_in_progress', 'Jobs currently processing')
processing_time = Histogram('worker_job_processing_seconds', 'Job processing time')

@processing_time.time()
def process_job(job):
    jobs_in_progress.inc()
    try:
        do_work(job)
        jobs_processed.labels(queue=job['queue_name'], status='completed').inc()
    except:
        jobs_processed.labels(queue=job['queue_name'], status='failed').inc()
        raise
    finally:
        jobs_in_progress.dec()

Next Steps