En los posts anteriores exploramos las goroutines, channels, select, context y las primitivas de sincronización. Hoy vamos a combinar todos estos conceptos para crear uno de los patrones más útiles en Go: los Worker Pools. Si vienes de Java, los worker pools son similares a ThreadPoolExecutor, pero mucho más simples y expresivos.
¿Qué es un Worker Pool?
Un Worker Pool es un patrón que limita el número de goroutines concurrentes que procesan tareas. En lugar de crear una goroutine por cada tarea (lo cual puede ser costoso), creas un número fijo de workers que procesan tareas de una cola.
¿Por Qué Usar Worker Pools?
Problema sin Worker Pool:
for i := 0; i < 10000; i++ {
go processTask(i)
}
Problemas:
❌ Demasiadas goroutines pueden consumir mucha memoria
❌ Overhead de scheduling
❌ Puede saturar recursos del sistema
❌ Difícil controlar el throughput
Solución con Worker Pool:
numWorkers := 10
jobs := make(chan Task, 100)
for w := 0; w < numWorkers; w++ {
go worker(jobs)
}
for i := 0; i < 10000; i++ {
jobs <- Task{ID: i}
}
Ventajas:
Comparación: Java vs Go
| Aspecto | Java (ThreadPoolExecutor) | Go (Worker Pool) |
| Configuración | Compleja (corePoolSize, maxPoolSize, queue) | Simple (número de workers) |
| Sintaxis | Verbosa | Simple y expresiva |
| Overhead | Alto (threads del OS) | Bajo (goroutines) |
| Escalabilidad | Limitada por threads | Millones de goroutines posibles |
| Cancelación | Compleja | Simple con context |
Worker Pool Básico
Empecemos con un ejemplo simple:
type Task struct {
ID int
Data string
}
type Result struct {
TaskID int
Output string
Error error
}
func basicWorkerPool() {
numWorkers := 3
numTasks := 10
jobs := make(chan Task, numTasks)
results := make(chan Result, numTasks)
var wg sync.WaitGroup
for w := 1; w <= numWorkers; w++ {
wg.Add(1)
go func(workerID int) {
defer wg.Done()
for task := range jobs {
time.Sleep(200 * time.Millisecond)
output := fmt.Sprintf("Processed by worker %d: %s", workerID, task.Data)
results <- Result{
TaskID: task.ID,
Output: output,
Error: nil,
}
}
}(w)
}
for i := 1; i <= numTasks; i++ {
jobs <- Task{
ID: i,
Data: fmt.Sprintf("Task %d", i),
}
}
close(jobs)
go func() {
wg.Wait()
close(results)
}()
for result := range results {
fmt.Printf("Result: %s\n", result.Output)
}
}
Componentes clave:
Channel de jobs: Cola de tareas pendientes
Channel de results: Resultados procesados
Workers: Goroutines que procesan tareas
WaitGroup: Esperar que todos los workers terminen
Flujo:
Crear N workers que leen de jobs
Enviar tareas a jobs
Cerrar jobs cuando no hay más tareas
Workers terminan cuando jobs se cierra
Cerrar results cuando todos los workers terminan
Recibir resultados de results
Worker Pool con Context (Cancelación)
Agregar cancelación con context es esencial para worker pools robustos:
func workerPoolWithContext() {
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()
numWorkers := 3
jobs := make(chan Task, 10)
results := make(chan Result, 10)
var wg sync.WaitGroup
for w := 1; w <= numWorkers; w++ {
wg.Add(1)
go func(workerID int) {
defer wg.Done()
for {
select {
case <-ctx.Done():
fmt.Printf("Worker %d: Cancelled\n", workerID)
return
case task, ok := <-jobs:
if !ok {
return
}
time.Sleep(500 * time.Millisecond)
select {
case results <- Result{
TaskID: task.ID,
Output: fmt.Sprintf("Processed %s", task.Data),
}:
case <-ctx.Done():
return
}
}
}
}(w)
}
go func() {
defer close(jobs)
for i := 1; i <= 10; i++ {
select {
case jobs <- Task{ID: i, Data: fmt.Sprintf("Task %d", i)}:
case <-ctx.Done():
return
}
}
}()
go func() {
wg.Wait()
close(results)
}()
for result := range results {
fmt.Printf("Result: %s\n", result.Output)
}
}
Características importantes:
✅ Workers verifican ctx.Done() antes y durante el procesamiento
✅ El sender de tareas también verifica cancelación
✅ El envío de resultados puede cancelarse
✅ Timeout automático después de 2 segundos
Worker Pool Reutilizable
Crear una estructura reutilizable hace el código más mantenible:
type WorkerPool struct {
numWorkers int
jobs chan Task
results chan Result
wg sync.WaitGroup
}
func NewWorkerPool(numWorkers, jobBufferSize int) *WorkerPool {
return &WorkerPool{
numWorkers: numWorkers,
jobs: make(chan Task, jobBufferSize),
results: make(chan Result, jobBufferSize),
}
}
func (wp *WorkerPool) Start(processor func(Task) Result) {
for i := 0; i < wp.numWorkers; i++ {
wp.wg.Add(1)
go func(workerID int) {
defer wp.wg.Done()
for task := range wp.jobs {
result := processor(task)
result.TaskID = task.ID
wp.results <- result
}
}(i)
}
}
func (wp *WorkerPool) Submit(task Task) {
wp.jobs <- task
}
func (wp *WorkerPool) Close() {
close(wp.jobs)
wp.wg.Wait()
close(wp.results)
}
func (wp *WorkerPool) Results() <-chan Result {
return wp.results
}
Uso:
pool := NewWorkerPool(3, 10)
processor := func(task Task) Result {
time.Sleep(200 * time.Millisecond)
return Result{
Output: fmt.Sprintf("Processed: %s", task.Data),
Error: nil,
}
}
pool.Start(processor)
for i := 1; i <= 5; i++ {
pool.Submit(Task{
ID: i,
Data: fmt.Sprintf("Task %d", i),
})
}
pool.Close()
for result := range pool.Results() {
fmt.Printf("Result: %s\n", result.Output)
}
Ventajas:
Ejemplo Práctico: Procesar Miles de Tareas
Un ejemplo real procesando 1000 tareas con 10 workers:
func processThousandsOfTasks() {
numWorkers := 10
numTasks := 1000
jobs := make(chan int, numTasks)
results := make(chan int, numTasks)
var wg sync.WaitGroup
startTime := time.Now()
for w := 0; w < numWorkers; w++ {
wg.Add(1)
go func() {
defer wg.Done()
for taskID := range jobs {
result := taskID * 2
time.Sleep(10 * time.Millisecond)
results <- result
}
}()
}
for i := 1; i <= numTasks; i++ {
jobs <- i
}
close(jobs)
go func() {
wg.Wait()
close(results)
}()
count := 0
for range results {
count++
if count%100 == 0 {
fmt.Printf("Processed %d tasks...\n", count)
}
}
elapsed := time.Since(startTime)
fmt.Printf("Processed %d tasks in %v\n", count, elapsed)
fmt.Printf("Throughput: %.2f tasks/second\n", float64(count)/elapsed.Seconds())
}
Resultado esperado:
Processed 100 tasks...
Processed 200 tasks...
...
Processed 1000 tasks in ~1s
Throughput: ~1000 tasks/second
Análisis:
Con 10 workers y 10ms por tarea: ~100 tareas/segundo teórico
El buffer permite enviar todas las tareas sin bloqueo
Los workers procesan en paralelo
Worker Pool con Prioridades
A veces necesitas procesar tareas con diferentes prioridades:
type PriorityTask struct {
Task
Priority int
}
func priorityWorkerPool() {
highPriority := make(chan PriorityTask, 10)
lowPriority := make(chan PriorityTask, 10)
var wg sync.WaitGroup
wg.Add(1)
go func() {
defer wg.Done()
for {
select {
case task := <-highPriority:
fmt.Printf("HIGH PRIORITY: Processing %s\n", task.Data)
time.Sleep(100 * time.Millisecond)
case task := <-lowPriority:
fmt.Printf("LOW PRIORITY: Processing %s\n", task.Data)
time.Sleep(100 * time.Millisecond)
default:
time.Sleep(50 * time.Millisecond)
if len(highPriority) == 0 && len(lowPriority) == 0 {
return
}
}
}
}()
lowPriority <- PriorityTask{Task: Task{ID: 1, Data: "Low 1"}, Priority: 1}
highPriority <- PriorityTask{Task: Task{ID: 2, Data: "High 1"}, Priority: 10}
lowPriority <- PriorityTask{Task: Task{ID: 3, Data: "Low 2"}, Priority: 1}
highPriority <- PriorityTask{Task: Task{ID: 4, Data: "High 2"}, Priority: 10}
time.Sleep(1 * time.Second)
wg.Wait()
}
Características:
✅ select procesa primero highPriority (se evalúa primero)
✅ Las tareas de alta prioridad se procesan antes
✅ Múltiples workers pueden procesar diferentes prioridades
Worker Pool con Rate Limiting
Limitar el rate de procesamiento puede ser útil:
func rateLimitedWorkerPool(rate time.Duration) {
numWorkers := 5
jobs := make(chan Task, 100)
results := make(chan Result, 100)
limiter := make(chan struct{}, numWorkers)
go func() {
ticker := time.NewTicker(rate)
defer ticker.Stop()
for range ticker.C {
select {
case limiter <- struct{}{}:
default:
}
}
}()
var wg sync.WaitGroup
for w := 0; w < numWorkers; w++ {
wg.Add(1)
go func() {
defer wg.Done()
for task := range jobs {
<-limiter
result := processTask(task)
results <- result
}
}()
}
}
Worker Pool con Error Handling
Manejar errores correctamente es crucial:
func workerPoolWithErrors() {
jobs := make(chan Task, 10)
results := make(chan Result, 10)
errors := make(chan error, 10)
var wg sync.WaitGroup
for w := 0; w < 3; w++ {
wg.Add(1)
go func() {
defer wg.Done()
for task := range jobs {
result, err := processTaskWithError(task)
if err != nil {
select {
case errors <- err:
default:
log.Printf("Error channel full: %v", err)
}
continue
}
results <- result
}
}()
}
go func() {
for err := range errors {
fmt.Printf("Error: %v\n", err)
}
}()
}
Mejores Prácticas
1. Elegir el Número Correcto de Workers
numWorkers := runtime.NumCPU()
numWorkers := runtime.NumCPU() * 2
numWorkers := getConfig().WorkerCount
numWorkers := 100
Regla general:
CPU-bound: runtime.NumCPU()
I/O-bound: runtime.NumCPU() * 2 o más
Network-bound: Puede ser mucho más alto
2. Usar Buffers Apropiados
jobs := make(chan Task, numTasks)
jobs := make(chan Task, numWorkers*2)
jobs := make(chan Task, 1)
jobs := make(chan Task, 1000000)
3. Siempre Cerrar Channels Correctamente
for i := 0; i < numTasks; i++ {
jobs <- Task{ID: i}
}
close(jobs)
go func() {
wg.Wait()
close(results)
}()
for i := 0; i < numTasks; i++ {
jobs <- Task{ID: i}
}
4. Usar Context para Cancelación
func workerPool(ctx context.Context, numWorkers int) {
select {
case <-ctx.Done():
return
case task := <-jobs:
processTask(task)
}
}
func workerPool(numWorkers int) {
}
5. Manejar Errores Correctamente
errors := make(chan error, 10)
go func() {
for err := range errors {
log.Printf("Error: %v", err)
}
}()
result, err := processTask(task)
if err != nil {
}
6. Monitorear el Pool
type WorkerPoolStats struct {
TotalTasks int64
CompletedTasks int64
FailedTasks int64
ActiveWorkers int
}
func (wp *WorkerPool) Stats() WorkerPoolStats {
return WorkerPoolStats{
TotalTasks: atomic.LoadInt64(&wp.totalTasks),
CompletedTasks: atomic.LoadInt64(&wp.completedTasks),
FailedTasks: atomic.LoadInt64(&wp.failedTasks),
ActiveWorkers: len(wp.jobs),
}
}
Errores Comunes
❌ Error 1: Olvidar Cerrar el Channel de Jobs
for i := 0; i < 10; i++ {
jobs <- Task{ID: i}
}
wg.Wait()
for i := 0; i < 10; i++ {
jobs <- Task{ID: i}
}
close(jobs)
wg.Wait()
❌ Error 2: Cerrar Results Demasiado Pronto
close(jobs)
close(results)
wg.Wait()
close(jobs)
go func() {
wg.Wait()
close(results)
}()
❌ Error 3: Demasiados o Muy Pocos Workers
numWorkers := 10000
numWorkers := 1
numWorkers := runtime.NumCPU()
numWorkers := runtime.NumCPU() * 2
❌ Error 4: No Verificar Cancelación en Workers
for task := range jobs {
processTask(task)
}
for {
select {
case <-ctx.Done():
return
case task, ok := <-jobs:
if !ok {
return
}
processTask(task)
}
}
❌ Error 5: Buffer Muy Pequeño o Muy Grande
jobs := make(chan Task, 1)
jobs := make(chan Task, 1000000)
jobs := make(chan Task, numWorkers*2)
Comparación: Java vs Go
Java: ThreadPoolExecutor
ExecutorService executor = Executors.newFixedThreadPool(10);
List<Future<String>> futures = new ArrayList<>();
for (int i = 0; i < 100; i++) {
final int taskId = i;
Future<String> future = executor.submit(() -> {
return processTask(taskId);
});
futures.add(future);
}
for (Future<String> future : futures) {
try {
String result = future.get();
System.out.println(result);
} catch (ExecutionException e) {
e.printStackTrace();
}
}
executor.shutdown();
Problemas:
Go: Worker Pool
numWorkers := 10
jobs := make(chan Task, 100)
results := make(chan Result, 100)
var wg sync.WaitGroup
for w := 0; w < numWorkers; w++ {
wg.Add(1)
go func() {
defer wg.Done()
for task := range jobs {
results <- processTask(task)
}
}()
}
for i := 0; i < 100; i++ {
jobs <- Task{ID: i}
}
close(jobs)
go func() {
wg.Wait()
close(results)
}()
for result := range results {
fmt.Println(result)
}
Ventajas:
Conclusiones
Los Worker Pools son un patrón esencial en Go:
✅ Control de concurrencia: Limita el número de goroutines activas
✅ Eficiencia: Mejor uso de recursos que crear una goroutine por tarea
✅ Escalabilidad: Puede procesar miles o millones de tareas
✅ Cancelación: Fácil agregar cancelación con context
✅ Flexibilidad: Puede adaptarse a diferentes necesidades (prioridades, rate limiting, etc.)
✅ Simplicidad: Código claro y mantenible
Si vienes de Java, los worker pools en Go son mucho más simples y expresivos. La combinación de goroutines, channels y WaitGroup hace que crear worker pools sea trivial comparado con ThreadPoolExecutor.
Próximos Pasos
En el siguiente post exploraremos patrones avanzados de concurrencia en Go, incluyendo pipelines, fan-out/fan-in, y otros patrones que combinan todos los conceptos que hemos aprendido.
¿Has usado worker pools en tus proyectos de Go? ¿Qué número de workers usas y cómo lo determinas? Comparte tus experiencias y casos de uso en los comentarios. Y si quieres ver el código completo de estos ejemplos, puedes encontrarlo en mi repositorio go-mastery-lab.