Go Worker Pool y procesamiento concurrente de tareas: patrones para sistemas de alta carga
Introducción: por qué necesitas un worker pool en servicios Go
En 2026, Go sigue siendo uno de los lenguajes dominantes para construir sistemas backend de alta carga. Las goroutines son económicas —puedes lanzar millones de ellas—, pero eso no significa que el spawn descontrolado de goroutines sea la estrategia correcta. En servicios reales, el paralelismo ilimitado conduce al agotamiento de descriptores de archivo, presión sobre el recolector de basura, sobrecarga de dependencias downstream (bases de datos, APIs de terceros) y latencias impredecibles.
El worker pool resuelve problemas clave:
- Limitación del paralelismo — no más de N tareas se ejecutan simultáneamente.
- Reutilización de goroutines — en lugar de crear una goroutine por cada solicitud, se usa un pool fijo.
- Backpressure — cuando la cola se desborda, el productor se bloquea o recibe un error, en lugar de generar carga de forma descontrolada.
- Aislamiento de errores — un panic en un worker no tumba todo el programa si existe un recover correcto.
En este artículo cubriremos todo de forma progresiva: desde una implementación primitiva hasta una solución production-ready con escalado dinámico, lógica de retry, graceful shutdown, integración con PostgreSQL como cola persistente y monitoreo completo.
Implementación básica de worker pool con goroutines y canales
El worker pool clásico en Go se construye sobre tres primitivos: un canal de tareas (jobs chan), un canal de resultados (results chan) y sync.WaitGroup para esperar a que todos los workers finalicen.
package workerpool
import (
"context"
"fmt"
"sync"
)
// Job representa una unidad de trabajo a ejecutar.
type Job struct {
ID int
Payload any
}
// Result almacena el resultado de la ejecución de una tarea.
type Result struct {
JobID int
Value any
Err error
}
// ProcessFunc — función de procesamiento de una tarea, se pasa al crear el pool.
type ProcessFunc func(ctx context.Context, job Job) Result
// Pool — worker pool básico.
type Pool struct {
jobs chan Job
results chan Result
wg sync.WaitGroup
process ProcessFunc
size int
}
// NewPool crea un pool con size workers y un buffer de cola queueSize.
func NewPool(size, queueSize int, fn ProcessFunc) *Pool {
return &Pool{
jobs: make(chan Job, queueSize),
results: make(chan Result, queueSize),
process: fn,
size: size,
}
}
// Start lanza los workers e inicia el procesamiento de tareas.
func (p *Pool) Start(ctx context.Context) {
for i := 0; i < p.size; i++ {
p.wg.Add(1)
go p.worker(ctx)
}
}
// worker — goroutine que lee tareas del canal jobs.
func (p *Pool) worker(ctx context.Context) {
defer p.wg.Done()
for {
select {
case job, ok := <-p.jobs:
if !ok {
// Canal cerrado — el worker finaliza.
return
}
result := p.process(ctx, job)
// Enviamos el resultado; si el canal está lleno — bloqueamos.
select {
case p.results <- result:
case <-ctx.Done():
return
}
case <-ctx.Done():
return
}
}
}
// Submit envía una tarea a la cola.
// Devuelve error si el contexto fue cancelado o el canal está lleno.
func (p *Pool) Submit(ctx context.Context, job Job) error {
select {
case p.jobs <- job:
return nil
case <-ctx.Done():
return fmt.Errorf("worker pool: submit cancelled: %w", ctx.Err())
}
}
// Results devuelve el canal de resultados para lectura por el consumidor.
func (p *Pool) Results() <-chan Result {
return p.results
}
// Stop cierra el canal de tareas y espera a que todos los workers finalicen.
func (p *Pool) Stop() {
close(p.jobs)
p.wg.Wait()
close(p.results)
}
Decisiones arquitectónicas clave en esta implementación:
- El canal
jobsestá bufferizado — el productor no se bloquea de inmediato ante picos de carga breves. - El worker escucha dos canales simultáneamente mediante
select: tareas y cancelación de contexto — esto evita bloqueos durante el shutdown. - Cerrar
jobses la señal a los workers para finalizar: Go garantiza que todos los valores enviados antes del cierre serán leídos.
Escalado dinámico del pool: adaptación a la carga
Un pool estático es óptimo con carga predecible. En sistemas reales, la carga varía: de noche el tráfico es mínimo, en horas pico puede ser 10 veces mayor. Un pool dinámico permite escalar el número de workers en el rango [minWorkers, maxWorkers] según una métrica: la longitud de la cola de tareas.
package workerpool
import (
"context"
"sync"
"sync/atomic"
"time"
)
// DynamicPool escala el número de workers según la carga.
type DynamicPool struct {
jobs chan Job
results chan Result
process ProcessFunc
mu sync.Mutex
wg sync.WaitGroup
ctx context.Context
cancel context.CancelFunc
active atomic.Int64 // número actual de workers
minWorkers int
maxWorkers int
}
func NewDynamicPool(min, max, queueSize int, fn ProcessFunc) *DynamicPool {
ctx, cancel := context.WithCancel(context.Background())
p := &DynamicPool{
jobs: make(chan Job, queueSize),
results: make(chan Result, queueSize),
process: fn,
ctx: ctx,
cancel: cancel,
minWorkers: min,
maxWorkers: max,
}
// Lanzamos el número mínimo de workers de inmediato.
for i := 0; i < min; i++ {
p.startWorker()
}
// Lanzamos la goroutine de autoescalado.
go p.autoscale()
return p
}
func (p *DynamicPool) startWorker() {
p.active.Add(1)
p.wg.Add(1)
go func() {
defer func() {
p.active.Add(-1)
p.wg.Done()
}()
for {
select {
case job, ok := <-p.jobs:
if !ok {
return
}
result := p.process(p.ctx, job)
select {
case p.results <- result:
case <-p.ctx.Done():
return
}
case <-p.ctx.Done():
return
}
}
}()
}
// autoscale verifica la carga cada 500 ms y añade workers si es necesario.
func (p *DynamicPool) autoscale() {
ticker := time.NewTicker(500 * time.Millisecond)
defer ticker.Stop()
for {
select {
case <-ticker.C:
queueLen := len(p.jobs)
current := int(p.active.Load())
// Si la cola está más del 50% llena y se puede crecer — añadimos un worker.
if queueLen > cap(p.jobs)/2 && current < p.maxWorkers {
p.startWorker()
}
// Si la cola está vacía y hay más workers que el mínimo — detenemos uno.
// Para detener un worker individual se usa un canal drain (simplificado).
case <-p.ctx.Done():
return
}
}
}
func (p *DynamicPool) Stop() {
p.cancel()
close(p.jobs)
p.wg.Wait()
close(p.results)
}
Para reducir el número de workers en implementaciones de producción se utiliza un canal quit individual por worker, o la librería golang.org/x/sync/semaphore para un control preciso del paralelismo.
Manejo de errores y lógica de retry sin fugas de goroutines
En sistemas de alta carga, parte de las tareas inevitablemente falla: errores de red, indisponibilidad temporal de dependencias, deadlines de contexto. Un retry ingenuo en bucle dentro del worker lo bloquea y reduce el throughput del pool. El enfoque correcto es reenviar la tarea a la cola con retardo exponencial y un número limitado de intentos.
package workerpool
import (
"context"
"errors"
"log/slog"
"math"
"time"
)
// RetryableJob extiende Job con metadatos para retry.
type RetryableJob struct {
Job
Attempt int
MaxRetries int
NextRunAt time.Time
}
// withRetry envuelve ProcessFunc añadiendo lógica de retry.
func withRetry(pool *Pool, fn ProcessFunc) ProcessFunc {
return func(ctx context.Context, job Job) Result {
rj, ok := job.Payload.(RetryableJob)
if !ok {
return fn(ctx, job)
}
result := fn(ctx, job)
if result.Err == nil {
return result
}
// No reintentamos si el contexto fue cancelado.
if errors.Is(result.Err, context.Canceled) || errors.Is(result.Err, context.DeadlineExceeded) {
return result
}
if rj.Attempt >= rj.MaxRetries {
slog.Error("job failed permanently",
"job_id", job.ID,
"attempts", rj.Attempt,
"err", result.Err,
)
return result
}
// Retardo exponencial: 2^attempt segundos, máximo 60 segundos.
delay := time.Duration(math.Min(math.Pow(2, float64(rj.Attempt)), 60)) * time.Second
rj.Attempt++
rj.NextRunAt = time.Now().Add(delay)
// Planificamos el reenvío en una goroutine separada.
// La goroutine terminará tras el delay o al cancelarse el contexto — sin fugas.
go func() {
select {
case <-time.After(delay):
retryJob := Job{ID: job.ID, Payload: rj}
if err := pool.Submit(ctx, retryJob); err != nil {
slog.Warn("failed to resubmit job", "job_id", job.ID, "err", err)
}
case <-ctx.Done():
slog.Info("retry cancelled due to context", "job_id", job.ID)
}
}()
return Result{JobID: job.ID, Err: nil} // no consideramos el error temporal como definitivo
}
}
Punto crítico: la goroutine para el retry diferido siempre termina — ya sea por el temporizador o al cancelarse el contexto. Esto elimina las fugas de goroutines (goroutine leak), uno de los problemas más peligrosos en servicios Go de larga ejecución.
Graceful Shutdown: finalización correcta de todas las tareas
Un servicio en producción debe finalizar correctamente al recibir SIGTERM o SIGINT: completar las tareas ya aceptadas, no aceptar nuevas y liberar recursos.
package main
import (
"context"
"log/slog"
"os"
"os/signal"
"syscall"
"time"
"github.com/yourorg/workerpool"
)
func main() {
// Contexto de la aplicación: se cancela al recibir una señal del SO.
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT)
defer stop()
pool := workerpool.NewPool(10, 100, processImage)
pool.Start(ctx)
// Goroutine productora: se detiene al cancelarse el contexto.
go func() {
for i := 0; ; i++ {
select {
case <-ctx.Done():
slog.Info("producer stopped")
return
default:
}
job := workerpool.Job{ID: i, Payload: fmt.Sprintf("image_%d.jpg", i)}
if err := pool.Submit(ctx, job); err != nil {
slog.Warn("submit failed", "err", err)
}
}
}()
// Leemos resultados en una goroutine separada.
go func() {
for result := range pool.Results() {
if result.Err != nil {
slog.Error("job error", "job_id", result.JobID, "err", result.Err)
}
}
}()
// Esperamos la señal de finalización.
<-ctx.Done()
slog.Info("shutdown signal received, draining pool...")
// Damos a los workers hasta 30 segundos para completar las tareas pendientes.
shutdownCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
done := make(chan struct{})
go func() {
pool.Stop() // cierra el canal jobs y espera el WaitGroup
close(done)
}()
select {
case <-done:
slog.Info("graceful shutdown complete")
case <-shutdownCtx.Done():
slog.Error("shutdown timeout exceeded, forcing exit")
os.Exit(1)
}
}
El patrón signal.NotifyContext apareció en Go 1.16 y es la forma idiomática de vincular el ciclo de vida de la aplicación con las señales del SO. El timeout de shutdown (aquí 30 segundos) previene una espera infinita ante tareas colgadas.
Integración con PostgreSQL: cola de tareas confiable
Para tareas que no se pueden perder al reiniciar el servicio (envío de emails, procesamiento de pagos), una cola en memoria no es suficiente — los datos no son persistentes. PostgreSQL con el patrón FOR UPDATE SKIP LOCKED resuelve este problema, convirtiendo una tabla en una cola distribuida confiable.
-- Esquema de la tabla de tareas
CREATE TABLE jobs (
id BIGSERIAL PRIMARY KEY,
type TEXT NOT NULL,
payload JSONB NOT NULL,
status TEXT NOT NULL DEFAULT 'pending', -- pending | running | done | failed
attempts INT NOT NULL DEFAULT 0,
max_retries INT NOT NULL DEFAULT 3,
run_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
CREATE INDEX idx_jobs_status_run_at ON jobs (status, run_at)
WHERE status IN ('pending', 'failed');
package pgqueue
import (
"context"
"encoding/json"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
type PGJob struct {
ID int64
Type string
Payload json.RawMessage
Attempts int
MaxRetries int
}
type Queue struct {
db *pgxpool.Pool
}
func NewQueue(db *pgxpool.Pool) *Queue {
return &Queue{db: db}
}
// Dequeue captura atómicamente una tarea de la cola.
// FOR UPDATE SKIP LOCKED garantiza que workers concurrentes no tomen la misma tarea dos veces.
func (q *Queue) Dequeue(ctx context.Context) (*PGJob, error) {
tx, err := q.db.Begin(ctx)
if err != nil {
return nil, err
}
defer tx.Rollback(ctx)
var job PGJob
err = tx.QueryRow(ctx, `
UPDATE jobs
SET status = 'running',
attempts = attempts + 1,
updated_at = NOW()
WHERE id = (
SELECT id FROM jobs
WHERE status IN ('pending', 'failed')
AND run_at <= NOW()
AND attempts < max_retries
ORDER BY run_at ASC
FOR UPDATE SKIP LOCKED
LIMIT 1
)
RETURNING id, type, payload, attempts, max_retries
`).Scan(&job.ID, &job.Type, &job.Payload, &job.Attempts, &job.MaxRetries)
if err != nil {
if err == pgx.ErrNoRows {
return nil, nil // cola vacía
}
return nil, err
}
return &job, tx.Commit(ctx)
}
// Complete marca la tarea como completada exitosamente.
func (q *Queue) Complete(ctx context.Context, jobID int64) error {
_, err := q.db.Exec(ctx,
`UPDATE jobs SET status = 'done', updated_at = NOW() WHERE id = $1`,
jobID,
)
return err
}
// Fail marca la tarea como fallida y planifica un retry con retardo exponencial.
func (q *Queue) Fail(ctx context.Context, jobID int64, attempts int) error {
delay := time.Duration(1<<attempts) * time.Second // 2^attempts segundos
if delay > 10*time.Minute {
delay = 10 * time.Minute
}
_, err := q.db.Exec(ctx, `
UPDATE jobs
SET status = 'failed',
run_at = NOW() + $1::interval,
updated_at = NOW()
WHERE id = $2
`, delay.String(), jobID)
return err
}
La propiedad clave de FOR UPDATE SKIP LOCKED: las filas bloqueadas por otras transacciones se omiten en lugar de provocar un bloqueo. Esto hace el patrón escalable — decenas de workers pueden consultar la misma tabla concurrentemente sin bloqueos mutuos.
Para el bucle de polling del worker, usa pausas adaptativas: si la cola está vacía — aumenta el intervalo de consulta (por ejemplo, de 100 ms a 5 segundos); cuando aparezcan tareas — vuelve al intervalo mínimo.
Monitoreo del worker pool: Prometheus y Grafana
Sin métricas, el worker pool es una caja negra. Los indicadores clave para monitorear son: tamaño de la cola, número de workers activos, tiempo de procesamiento de tareas, tasa de errores y retries.
package metrics
import (
"time"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/promauto"
)
var (
// Longitud actual de la cola de tareas.
QueueDepth = promauto.NewGauge(prometheus.GaugeOpts{
Name: "worker_pool_queue_depth",
Help: "Number of jobs waiting in the queue.",
})
// Número de workers activos.
ActiveWorkers = promauto.NewGauge(prometheus.GaugeOpts{
Name: "worker_pool_active_workers",
Help: "Number of currently active workers.",
})
// Histograma del tiempo de procesamiento de una tarea.
JobDuration = promauto.NewHistogramVec(
prometheus.HistogramOpts{
Name: "worker_pool_job_duration_seconds",
Help: "Time spent processing a job.",
Buckets: prometheus.ExponentialBuckets(0.001, 2, 15), // 1ms..16s
},
[]string{"job_type", "status"}, // status: success | error
)
// Contador de tareas procesadas.
JobsTotal = promauto.NewCounterVec(
prometheus.CounterOpts{
Name: "worker_pool_jobs_total",
Help: "Total number of processed jobs.",
},
[]string{"job_type", "status"},
)
// Contador de retries.
RetriesTotal = promauto.NewCounter(prometheus.CounterOpts{
Name: "worker_pool_retries_total",
Help: "Total number of job retries.",
})
)
// RecordJobExecution instrumenta la ejecución de una tarea.
func RecordJobExecution(jobType string, start time.Time, err error) {
duration := time.Since(start).Seconds()
status := "success"
if err != nil {
status = "error"
}
JobDuration.WithLabelValues(jobType, status).Observe(duration)
JobsTotal.WithLabelValues(jobType, status).Inc()
}
Alertas recomendadas en Grafana:
worker_pool_queue_depth > 1000durante 5 minutos — señal de sobrecarga o workers bloqueados.rate(worker_pool_jobs_total{status="error"}[5m]) / rate(worker_pool_jobs_total[5m]) > 0.05— tasa de errores superior al 5%.histogram_quantile(0.99, worker_pool_job_duration_seconds) > 30— latencia p99 superior a 30 segundos.
Comparación con soluciones existentes: cuándo implementar la propia
Antes de implementar tu propio worker pool, conviene evaluar las soluciones disponibles:
- asynq — cola basada en Redis con UI enriquecida (Asynq Monitor), soporte de schedules, prioridades y tareas únicas. Excelente opción si Redis ya está en el stack y se necesita integración rápida.
- river — cola de tareas nativa de PostgreSQL de los autores de pgx. Usa
LISTEN/NOTIFYen lugar de polling y soporta enqueue transaccional (la tarea se añade atómicamente junto con la transacción de negocio). Ideal si PostgreSQL es la única dependencia. - machinery — solución más pesada con soporte de múltiples brokers (Redis, AMQP, SQS), adecuada para workflows distribuidos.
Cuándo tiene sentido implementar tu propio worker pool:
- Las tareas son in-memory y no se requiere persistencia (por ejemplo, pipeline de procesamiento de solicitudes).
- Se necesita lógica específica de priorización o enrutamiento de tareas.
- Minimizar dependencias externas es crítico (servicios embebidos, edge).
- Las soluciones existentes son excesivas en funcionalidad y añaden complejidad indeseable.
Ejemplo práctico: servicio de procesamiento asíncrono de imágenes
Reunamos todo en un ejemplo realista: un servicio HTTP acepta cargas de imágenes, coloca tareas de procesamiento (redimensionado, conversión de formatos) en una cola PostgreSQL, y el worker pool toma las tareas y las procesa concurrentemente.
package main
import (
"context"
"encoding/json"
"log/slog"
"net/http"
"os"
"os/signal"
"syscall"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/prometheus/client_golang/prometheus/promhttp"
"github.com/yourorg/metrics"
"github.com/yourorg/pgqueue"
)
type ImagePayload struct {
S3Key string `json:"s3_key"`
Widths []int `json:"widths"`
Formats []string `json:"formats"`
}
// processImageJob — lógica de negocio para procesar una imagen.
func processImageJob(ctx context.Context, payload json.RawMessage) error {
var p ImagePayload
if err := json.Unmarshal(payload, &p); err != nil {
return err
}
// Aquí: descargar de S3, redimensionar, guardar de vuelta.
// Para el ejemplo emulamos el trabajo.
slog.Info("processing image", "s3_key", p.S3Key, "widths", p.Widths)
time.Sleep(100 * time.Millisecond) // simulación de trabajo
return nil
}
func runWorkers(ctx context.Context, queue *pgqueue.Queue, workerCount int) {
sem := make(chan struct{}, workerCount) // semáforo para el número de workers
for {
select {
case <-ctx.Done():
// Esperamos a que se liberen todos los slots del semáforo.
for i := 0; i < workerCount; i++ {
sem <- struct{}{}
}
return
default:
}
job, err := queue.Dequeue(ctx)
if err != nil {
slog.Error("dequeue error", "err", err)
time.Sleep(time.Second)
continue
}
if job == nil {
// Cola vacía — pausa adaptativa.
time.Sleep(200 * time.Millisecond)
continue
}
// Ocupamos un slot del semáforo.
sem <- struct{}{}
metrics.ActiveWorkers.Inc()
go func(j *pgqueue.PGJob) {
defer func() {
<-sem // liberamos el slot
metrics.ActiveWorkers.Dec()
}()
start := time.Now()
err := processImageJob(ctx, j.Payload)
metrics.RecordJobExecution(j.Type, start, err)
if err != nil {
slog.Error("job failed", "job_id", j.ID, "err", err)
if qErr := queue.Fail(ctx, j.ID, j.Attempts); qErr != nil {
slog.Error("failed to mark job as failed", "err", qErr)
}
return
}
if err := queue.Complete(ctx, j.ID); err != nil {
slog.Error("failed to mark job complete", "err", err)
}
}(job)
}
}
func main() {
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT)
defer stop()
db, err := pgxpool.New(ctx, os.Getenv("DATABASE_URL"))
if err != nil {
slog.Error("db connect failed", "err", err)
os.Exit(1)
}
defer db.Close()
queue := pgqueue.NewQueue(db)
// Servidor HTTP: acepta cargas y expone métricas.
mux := http.NewServeMux()
mux.Handle("/metrics", promhttp.Handler())
mux.HandleFunc("/upload", func(w http.ResponseWriter, r *http.Request) {
// Simplificado: codificamos el payload y lo insertamos en la cola.
payload, _ := json.Marshal(ImagePayload{
S3Key: "uploads/" + r.URL.Query().Get("file"),
Widths: []int{320, 640, 1280},
Formats: []string{"webp", "jpg"},
})
if _, err := db.Exec(r.Context(),
`INSERT INTO jobs (type, payload) VALUES ('image_process', $1)`, payload); err != nil {
http.Error(w, "enqueue failed", http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusAccepted)
})
srv := &http.Server{Addr: ":8080", Handler: mux}
go srv.ListenAndServe()
// Lanzamos el pool de workers.
go runWorkers(ctx, queue, 20)
<-ctx.Done()
slog.Info("shutting down")
shutCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
srv.Shutdown(shutCtx)
}
Este ejemplo demuestra el ciclo completo: HTTP → cola PostgreSQL → worker pool con semáforo → métricas Prometheus → graceful shutdown. El semáforo (chan struct{}) reemplaza aquí un pool explícito de goroutines y es el patrón idiomático de Go para limitar el paralelismo.
Conclusión
El worker pool en Go no es solo un patrón de optimización, sino una herramienta fundamental para construir servicios confiables de alta carga. Hemos cubierto el espectro completo: desde la implementación básica con canales hasta el escalado dinámico, retry sin fugas de goroutines, graceful shutdown y PostgreSQL como cola persistente con FOR UPDATE SKIP LOCKED.
Conclusiones clave:
- Limita siempre el paralelismo de forma explícita — mediante el tamaño del pool, un semáforo o un canal bufferizado.
- El contexto debe atravesar todo el stack: desde el handler HTTP hasta la consulta a la base de datos.
- El graceful shutdown con timeout es obligatorio en producción — el SO no espera indefinidamente.
- Para tareas persistentes, PostgreSQL con
FOR UPDATE SKIP LOCKEDes una solución madura y confiable, especialmente en conjunto con la librería river. - Las métricas de Prometheus deben integrarse desde el primer día — no las añadas de forma retroactiva.
Go proporciona todos los primitivos necesarios para implementar un worker pool production-ready sin dependencias externas. Comprender estos patrones distingue a un desarrollador capaz de construir sistemas resilientes ante cargas reales.
Tecnologías
Etiquetas
Ruslan Ismailov
Desarrollador Senior Web / Backend. Desarrollador senior web/backend con 9 años de experiencia. Stack: PHP, Laravel, PostgreSQL, Redis, Docker, Kubernetes, REST, microservicios, CI/CD. Más sobre mí →