Inrtoduction
After covering Go's concurrency primitives in previous articles, it's time to apply those concepts through common concurrency patterns. Among the most useful patterns in Go development:
- Worker-Pool Pattern
- Producer-Consumer Pattern
- Pipeline Pattern
- Event-Driven Pattern
- Reactor Pattern
- Futures and Promises Pattern
Let's start with the Worker-Pool Pattern, which controls concurrent task execution, improves system throughput, and manages system resources more efficiently.
This article demonstrates how to build a Worker Pool implemantation called GoPool, incrementally improving it with AI assistance. We'll compare our implementation against popular open-source alternatives in terms of feature completeness, usability, and performance.
Design Phase
Defining Core Requirements
A production-ready Worker Pool typically needs the following capabilities:
- Task Queue: Thread-safe queue for pending tasks accessible by multiple goroutines
- Worker Threads: Goroutines that consume tasks from the queue and execute them
- Pool Size Control: Mechanism to limit concurrent workers and prevent resource exhaustion
- Graceful Shutdown: Ensure all in-progress tasks complete before terminating the pool
- Error Handling: Capture and report task execution errors
- Timeout Control: Cancel long-running tasks that exceed defined thresholds
- Task Priority: Execute higher-priority tasks before lower-priority ones
- Result Retrieval: Return execution results for tasks that produce output
Implementing all these features requires careful planning. Rather than asking the AI to generate everything at once, we'll build incrementally—using AI as a programming assistant that generates specific components.
Project Setup
We'll create a repository named gopool on GitHub. The core entities are Task, Worker, and Pool. The Pool maintains a worker stack (instead of a queue—stacks are more flexible for implementing worker expiration logic). The Task is represented as a simple function type.
Initial skeleton code (~50 lines) establishes the foundation:
pool.go
type Pool struct {
capacity int
running int
workers WorkerStack
lock sync.Mutex
}
func (p *Pool) Submit(task Task)
func (p *Pool) Running() int
func (p *Pool) Cap() int
func (p *Pool) Release()
worker.go
type Task func()
type Worker struct {
pool *Pool
task chan Task
lastUsedTime time.Time
}
type WorkerStack interface {
Push(w *Worker)
Pop() *Worker
Len() int
}
After defining these structures, an IDE's autocomplete feature helped complete property declarations and comments.
Implementation
Basic Implementation
Initial implementation provides the core Submit() and Release() methods. The Submit method checks if there are available workers; if so, it pops one from the stack and assigns the task. When the pool reaches capacity, a sync.Cond condition variable handles blocking until a worker becomes available.
pool.go
package gopool
import (
"sync"
)
type Pool struct {
capacity int
running int
workers WorkerStack
lock sync.Mutex
cond *sync.Cond
}
func NewPool(capacity int) *Pool {
p := &Pool{
capacity: capacity,
workers: &workerStack{},
}
p.cond = sync.NewCond(&p.lock)
return p
}
func (p *Pool) Submit(task Task) {
p.lock.Lock()
defer p.lock.Unlock()
for p.running == p.capacity {
p.cond.Wait()
}
worker := p.workers.Pop()
if worker == nil {
worker = &Worker{
pool: p,
task: make(chan Task, 1),
}
}
worker.task <- task
worker.run()
p.running++
}
func (p *Pool) Running() int {
return p.running
}
func (p *Pool) Cap() int {
return p.capacity
}
func (p *Pool) Release() {
p.lock.Lock()
defer p.lock.Unlock()
for p.workers.Len() > 0 {
worker := p.workers.Pop()
close(worker.task)
}
p.running = 0
}
worker.go
package gopool
import (
"time"
)
type Task func()
type Worker struct {
pool *Pool
task chan Task
lastUsedTime time.Time
}
type WorkerStack interface {
Push(w *Worker)
Pop() *Worker
Len() int
}
type workerStack struct {
workers []*Worker
}
func (ws *workerStack) Push(w *Worker) {
ws.workers = append(ws.workers, w)
}
func (ws *workerStack) Pop() *Worker {
if len(ws.workers) == 0 {
return nil
}
w := ws.workers[len(ws.workers)-1]
ws.workers = ws.workers[:len(ws.workers)-1]
return w
}
func (ws *workerStack) Len() int {
return len(ws.workers)
}
func (w *Worker) run() {
go func() {
for task := range w.task {
if task == nil {
return
}
task()
w.pool.lock.Lock()
w.pool.workers.Push(w)
w.pool.running--
w.pool.lock.Unlock()
}
}()
}
Handling Pool Saturation
The sync.Cond mechanism enables graceful handling when the pool reaches capacity. When Submit() detects running == capacity, it calls Wait() to block the calling goroutine. When a Worker completes its task and returns to the pool, the run() method signals the condition variable to wake waiting goroutines.
Unit Testing
Basic test coverage validates the submit functionality:
package gopool
import (
"sync"
"testing"
)
func TestSubmit(t *testing.T) {
var wg sync.WaitGroup
p := NewPool(10)
for i := 0; i < 20; i++ {
wg.Add(1)
taskNum := i
task := func() {
t.Logf("Executing task %d", taskNum)
defer wg.Done()
}
p.Submit(task)
}
wg.Wait()
if p.Running() != 0 {
t.Errorf("Expected running workers to be 0, but got %d", p.Running())
}
}
A critical bug discovered during testing: the Worker channel was created without buffering (make(chan Task)), causing immediate blocking. The fix changed it to a buffered channel: make(chan Task, 1).
Test execution confirms correct behavior:
$ go test . -v === RUN TestSubmit pool_test.go:16: Executing task 9 pool_test.go:16: Executing task 7 pool_test.go:16: Executing task 8 ... --- PASS: TestSubmit (0.00s) PASS
<h2>Next Steps</h2>
<p>The current implementation (~240 lines total: ~200 from AI, ~20 manually written, ~20 from autocomplete) provides a functional Worker Pool with basic features. Future enhancements include:</p>
- Graceful shutdown with timeout support
- Task priority queues
- Configurable task timeout handling
- Performance optimization and benchmarking
- Context propagation to workers
Contributions and feedback welcome via GitHub Issues.