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