Portal/Notes πŸ“
Interview prep

Worker Pool β€” Go Concurrency Pattern

Implementasi worker pool pattern di Go: fan-out/fan-in, graceful shutdown, backpressure handling, context cancellation. Untuk batch processing, job queue, dan concurrent I/O.

Worker Pool β€” Concurrency Pattern

Use Case

Worker pool cocok untuk:

  • Batch processing: proses 10K record dari database
  • Concurrent I/O: fetch 100 API endpoints bersamaan
  • Job queue: process jobs dari message queue (RabbitMQ, Kafka)
  • CPU-bound work: image processing, encryption, hashing

Basic Implementation

package workerpool

import (
    "context"
    "fmt"
    "sync"
)

type Job struct {
    ID   int
    Data interface{}
}

type Result struct {
    Job    Job
    Output interface{}
    Err    error
}

type WorkerPool struct {
    workers int
    jobs    chan Job
    results chan Result
    wg      sync.WaitGroup
}

func New(workers int, queueSize int) *WorkerPool {
    return &WorkerPool{
        workers: workers,
        jobs:    make(chan Job, queueSize),
        results: make(chan Result, queueSize),
    }
}

// Start launches worker goroutines.
func (wp *WorkerPool) Start(ctx context.Context, handler func(Job) (interface{}, error)) {
    for i := 0; i < wp.workers; i++ {
        wp.wg.Add(1)
        go func(workerID int) {
            defer wp.wg.Done()
            for {
                select {
                case <-ctx.Done():
                    return
                case job, ok := <-wp.jobs:
                    if !ok {
                        return // channel closed
                    }
                    output, err := handler(job)
                    select {
                    case wp.results <- Result{Job: job, Output: output, Err: err}:
                    case <-ctx.Done():
                        return
                    }
                }
            }
        }(i)
    }
}

// Submit sends a job to the pool.
func (wp *WorkerPool) Submit(job Job) error {
    select {
    case wp.jobs <- job:
        return nil
    default:
        return fmt.Errorf("job queue full (size=%d)", cap(wp.jobs))
    }
}

// Results returns the results channel.
func (wp *WorkerPool) Results() <-chan Result {
    return wp.results
}

// Shutdown gracefully stops the pool.
func (wp *WorkerPool) Shutdown() {
    close(wp.jobs)    // stop accepting new jobs
    wp.wg.Wait()      // wait for workers to finish
    close(wp.results) // close results channel
}

Usage Example β€” Batch API Calls

func main() {
    ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
    defer cancel()

    pool := New(10, 100) // 10 workers, 100 job queue

    // Start workers
    pool.Start(ctx, func(job Job) (interface{}, error) {
        url := job.Data.(string)
        resp, err := http.Get(url)
        if err != nil {
            return nil, err
        }
        defer resp.Body.Close()
        return io.ReadAll(resp.Body)
    })

    // Submit jobs in background
    go func() {
        urls := []string{"https://api1.com", "https://api2.com", /* ... */}
        for i, url := range urls {
            if err := pool.Submit(Job{ID: i, Data: url}); err != nil {
                log.Printf("failed to submit job %d: %v", i, err)
            }
        }
        pool.Shutdown()
    }()

    // Collect results
    for result := range pool.Results() {
        if result.Err != nil {
            log.Printf("job %d failed: %v", result.Job.ID, result.Err)
            continue
        }
        log.Printf("job %d got %d bytes", result.Job.ID, len(result.Output.([]byte)))
    }
}

Backpressure Handling

Jangan cuma reject kalau queue penuh. Beri feedback ke upstream:

func (wp *WorkerPool) SubmitWithBackpressure(ctx context.Context, job Job) error {
    for {
        select {
        case <-ctx.Done():
            return ctx.Err()
        case wp.jobs <- job:
            return nil
        default:
            // Queue full β€” exponential backoff
            backoff := time.Millisecond * 10
            select {
            case <-ctx.Done():
                return ctx.Err()
            case <-time.After(backoff):
                backoff *= 2
                if backoff > time.Second {
                    backoff = time.Second
                }
            }
        }
    }
}

Semaphore Pattern (Tanpa Channel Pool)

Untuk kasus simpel, cukup pakai buffered channel sebagai semaphore:

func ProcessWithSemaphore(ctx context.Context, items []Item, maxConcurrency int) error {
    sem := make(chan struct{}, maxConcurrency)
    errs := make(chan error, len(items))
    var wg sync.WaitGroup

    for _, item := range items {
        select {
        case <-ctx.Done():
            return ctx.Err()
        case sem <- struct{}{}: // acquire
        }

        wg.Add(1)
        go func(it Item) {
            defer wg.Done()
            defer func() { <-sem }() // release

            if err := processOne(ctx, it); err != nil {
                errs <- err
            }
        }(item)
    }

    wg.Wait()
    close(errs)

    // Collect first error
    for err := range errs {
        if err != nil {
            return err
        }
    }
    return nil
}

Graceful Shutdown

func (wp *WorkerPool) GracefulShutdown(ctx context.Context, timeout time.Duration) error {
    // 1. Stop accepting new jobs
    close(wp.jobs)

    // 2. Wait for in-flight jobs with timeout
    done := make(chan struct{})
    go func() {
        wp.wg.Wait()
        close(done)
    }()

    select {
    case <-done:
        // All workers finished
        close(wp.results)
        return nil
    case <-time.After(timeout):
        // Timeout β€” cancel context to force shutdown
        return fmt.Errorf("shutdown timeout: %d workers still running", wp.ActiveWorkers())
    }
}

func (wp *WorkerPool) ActiveWorkers() int {
    // Track active workers with atomic counter in Start()
    return wp.activeWorkers.Load()
}

Panic Recovery per Worker

func (wp *WorkerPool) Start(ctx context.Context, handler func(Job) (interface{}, error)) {
    for i := 0; i < wp.workers; i++ {
        wp.wg.Add(1)
        go func(workerID int) {
            defer wp.wg.Done()
            defer func() {
                if r := recover(); r != nil {
                    log.Printf("WORKER %d PANIC: %v\nstack: %s", workerID, r, debug.Stack())
                    // Re-spawn worker
                    wp.wg.Add(1)
                    go func() { /* re-spawn */ }()
                }
            }()

            for job := range wp.jobs {
                output, err := handler(job)
                wp.results <- Result{Job: job, Output: output, Err: err}
            }
        }(i)
    }
}

Tuning: Berapa Jumlah Worker?

Rule of thumb:
- I/O-bound: workers = 2 Γ— CPU cores (atau 50-100 untuk network I/O)
- CPU-bound: workers = CPU cores (atau CPU cores - 1)
- Mixed: start dengan CPU cores Γ— 2, benchmark, adjust

Benchmark:
func BenchmarkWorkerPool(b *testing.B) {
    for workers := 1; workers <= 32; workers *= 2 {
        b.Run(fmt.Sprintf("workers=%d", workers), func(b *testing.B) {
            pool := New(workers, 1000)
            // ...
        })
    }
}

Interview Talking Points

"Kapan pakai worker pool vs semaphore?"

  • Worker pool: job yang banyak, long-running, butuh antrian + backpressure
  • Semaphore: task simpel, gak butuh job queue, cuma butuh concurrency limit
  • Rule: kalau task bisa selesai < 1 detik β†’ semaphore cukup. > 1 detik β†’ worker pool.

"Apa bedanya sama errgroup?"

  • golang.org/x/sync/errgroup: simpler, auto-cancel on first error
  • Worker pool: lebih fleksibel, bisa collect multiple errors, bisa pause/resume
  • Pilih errgroup untuk task independent yang gagal = semua gagal

"Bagaimana kalau worker crash?"

  • Panic recovery dengan defer recover() β€” log error, jangan crash seluruh service
  • Re-spawn worker kalau mati
  • Kalau worker mati berulang kali β†’ circuit breaker β†’ stop semua worker
Edit on GitHub

Last updated on

Worker Pool β€” Go Concurrency Pattern | Faisal Affan