Esta guía reúne ejemplos y contenido de varias fuentes. No todo el texto ni todo el código son originales.
Introducción
Esta guía presenta una base para implementar patrones habituales de concurrencia en Go, incluida la concurrencia limitada. Este patrón establece límites estructurales al trabajo simultáneo para proteger recursos finitos como la memoria y la capacidad de cómputo. Aunque el tema principal es la concurrencia limitada, aplicarla requiere entender el modelo de concurrencia de Go. La primera mitad explica sus mecanismos básicos: goroutines, canales, búferes y contextos. Si ya los conoces, puedes pasar a la sección sobre concurrencia limitada.
Conceptos
Goroutines
Una goroutine es una unidad de ejecución ligera gestionada por el entorno de ejecución de Go, en lugar de por el sistema operativo. No equivale a un hilo completo del sistema. Al invocar una función con go, se inicia una goroutine con una pila dinámica que puede crecer o reducirse. Su coste es bajo y pueden coexistir decenas de miles.
Fugas de goroutines
Una fuga de goroutines ocurre cuando una goroutine se inicia pero nunca termina. El recolector de basura no elimina las goroutines activas, por lo que estas se acumulan durante la vida del proceso y consumen memoria y otros recursos. Una causa habitual es una goroutine bloqueada al recibir de un canal al que nunca llega un valor:
func doWork(ch <-chan Job) {
go func() {
job := <-ch // blocks forever if nothing sends on ch and it is never closed
process(job)
}()
}
Si quien llama descarta ch sin cerrarlo ni enviar un valor, la goroutine queda bloqueada. Tras suficientes llamadas a doWork, el proceso acumula goroutines detenidas.
La solución consiste en proporcionar una vía de salida mediante una señal de cancelación de context, comprobada con select:
func doWork(ctx context.Context, ch <-chan Job) {
go func() {
select {
case job := <-ch:
process(job)
case <-ctx.Done(): // guaranteed exit when context is cancelled
return
}
}()
}
Esta es una de las razones fundamentales para propagar context.Context por el código concurrente. Además de los tiempos de espera, ofrece una señal de salida que las goroutines deben atender. Desde el lado de quien llama, defer wg.Done() permite que wg.Wait() no quede pendiente de una goroutine que terminó antes de tiempo.
Las fugas pueden detectarse durante las pruebas con goleak, que comprueba que no queden goroutines inesperadas al finalizar una prueba.
- Uber Go: goleak — Goroutine Leak Detector
- Dave Cheney: Never start a goroutine without knowing how it will stop
Canales
Un canal conecta goroutines concurrentes y permite enviar y recibir valores entre ellas. Transferir datos mediante canales puede hacer explícito el cambio de propiedad y ayudar a evitar condiciones de carrera. Es una alternativa a coordinar memoria compartida con sync.Mutex: compartir memoria mediante la comunicación.
Un sync.Mutex, o bloqueo de exclusión mutua, es una alternativa de menor nivel. Protege una variable compartida permitiendo que solo una goroutine mantenga el bloqueo a la vez. Las demás esperan hasta que quien lo posee llame a Unlock.
var mu sync.Mutex
var counter int
mu.Lock()
counter++ // only one goroutine can be here at a time
mu.Unlock()
Los mutex son adecuados cuando varias goroutines necesitan leer y modificar una estructura compartida. Los canales resultan útiles cuando se transfieren datos entre goroutines y conviene expresar claramente su propiedad. En la práctica, los proyectos de Go suelen utilizar ambos según el patrón de coordinación.
Canales sin búfer
Un canal sin búfer es síncrono y no tiene almacenamiento interno. El envío espera hasta que alguien pueda recibir el valor y la recepción espera hasta que alguien pueda enviarlo. Esta entrega directa permite sincronizar tareas: antes de continuar, el emisor sabe que el receptor ha recibido el trabajo. Es como entregar el testigo en una carrera de relevos: quien lo entrega no puede soltarlo hasta que la otra persona lo sujete. Ambas partes deben completar el intercambio.
Canales con búfer
Un canal con búfer incorpora almacenamiento de capacidad fija. Se crea, por ejemplo, con make(chan int, 5). Los emisores pueden enviar sin esperar mientras haya espacio; los receptores reciben a su ritmo y esperan cuando el búfer está vacío. Esto desacopla los tiempos de productores y consumidores y ayuda a absorber ráfagas de trabajo.
Es como una bandeja de entrada sobre una mesa: alguien puede dejar un documento y marcharse sin esperar a que lo leas. Si la bandeja se llena, la siguiente persona debe esperar a que quede espacio.
Contexto
El paquete context gestiona ciclos de vida, cancelaciones y plazos en cadenas de goroutines. Un Context se pasa como primer argumento a las funciones que deben participar en esa coordinación. Cuando se cancela una petición o vence un plazo, se cierra ctx.Done(). Las goroutines que atienden esa señal pueden detenerse y liberar sus recursos.
Referencia de sintaxis
Los patrones de concurrencia utilizan algunas construcciones específicas de Go que pueden resultar nuevas si vienes de otros lenguajes. Esta sección las explica brevemente.
El operador <-
<- es el operador de canales. Su posición respecto al nombre del canal indica si se envía o se recibe.
ch <- value // send: push value into ch; blocks if ch is full
value := <-ch // receive: pull a value out of ch; blocks if ch is empty
<-ch // receive and discard: unblocks when ch has a value, ignores it
La flecha refleja el flujo de datos: ch <- value envía un valor al canal y <-ch recibe un valor del canal.
En <-ctx.Done(), la recepción espera al canal de finalización. Como ese canal se cierra sin recibir envíos, la operación deja de bloquear cuando el contexto se cancela o vence su plazo. Es la forma habitual de esperar una señal de terminación.
make(chan T, n)
make inicializa canales, segmentos y mapas. En el caso de los canales:
ch := make(chan int) // unbuffered: zero capacity, synchronous handoff
ch := make(chan int, 5) // buffered: capacity 5, can hold 5 values before blocking
El segundo argumento indica la capacidad del búfer. Si se omite, la capacidad es cero y el canal no tiene búfer.
chan struct{} se utiliza habitualmente para señales sin datos. struct{} es una estructura vacía que ocupa cero bytes; lo relevante es la operación de envío o recepción.
select
select se parece a switch, pero selecciona operaciones sobre canales. Espera hasta que algún caso pueda continuar y lo ejecuta. Si varios están listos, Go elige uno de forma pseudoaleatoria.
select {
case job := <-jobs: // fires when jobs has a value to receive
handle(job)
case <-ctx.Done(): // fires when the context is cancelled
return
}
Un caso default hace que select no bloquee: si ningún otro caso está listo, se ejecuta inmediatamente.
select {
case job := <-jobs:
handle(job)
default:
// nothing ready, move on immediately
}
select permite esperar varios canales sin comprometerse con uno solo.
for range sobre un canal
for job := range ch recibe valores de un canal en un bucle. Entre iteraciones espera al siguiente valor y termina cuando el canal se ha cerrado y vaciado.
for job := range work {
process(job) // blocks here until work has a value or is closed
}
Este patrón permite que una goroutine trabajadora consuma tareas de un canal. El productor indica que ha terminado mediante close(work); el bucle procesa los valores pendientes antes de salir.
close
close(ch) cierra un canal. No se pueden enviar más valores: intentarlo provoca un pánico. Los receptores aún pueden consumir el búfer. Una vez vacío, las recepciones devuelven inmediatamente el valor cero.
El cierre actúa como una señal de difusión: desbloquea las recepciones que esperan ese canal. Así puede indicar un productor a sus trabajadores que ya no habrá más tareas.
defer
defer programa una llamada para cuando termine la función que la contiene, ya sea mediante una salida normal, un retorno anticipado o un pánico.
defer wg.Done() // called when the goroutine function exits
defer func() { <-sem }() // anonymous function deferred: releases semaphore token on exit
En código concurrente, defer asegura la limpieza: liberar un permiso de semáforo, decrementar un WaitGroup o cerrar un recurso aunque se produzca un retorno anticipado por error.
sync.WaitGroup
Un WaitGroup mantiene un contador de tareas. Quien llama lo incrementa antes de iniciar cada goroutine y espera a que llegue a cero.
var wg sync.WaitGroup
wg.Add(1) // increment before spawning
go func() {
defer wg.Done() // decrement when goroutine exits
doWork()
}()
wg.Wait() // blocks until counter reaches zero
Llama a wg.Add(1) antes de iniciar la goroutine. Si se hace dentro de ella, wg.Wait() podría observar un contador cero y regresar antes de que el trabajo se haya registrado.
context.WithTimeout y context.WithCancel
context.WithCancel(parent) devuelve un contexto derivado y una función cancel. Al llamar a cancel(), se cierra su canal Done y se notifica la cancelación a las goroutines que utilizan ese contexto.
ctx, cancel := context.WithCancel(context.Background())
defer cancel() // always call cancel to release resources
go func() {
select {
case <-ctx.Done(): // unblocks when cancel() is called
return
case job := <-jobs:
handle(job)
}
}()
context.WithTimeout(parent, duration) también cancela automáticamente cuando transcurre el plazo. Aunque el trabajo termine antes, conviene llamar a la función de cancelación mediante defer para liberar recursos.
ctx, cancel := context.WithTimeout(parentCtx, 5*time.Second)
defer cancel()
conn, err := dialer.DialContext(ctx, "tcp", addr) // fails if 5s elapses
go func()(...)()
Las goroutines anónimas suelen seguir este patrón:
go func(j Job) {
process(j)
}(job)
func(j Job) { ... } declara una función anónima con el parámetro j. El (job) final la invoca con el valor actual de job, y el prefijo go ejecuta esa llamada en una nueva goroutine.
Pasar job explícitamente deja claro qué valor usa cada goroutine y evita depender de una variable compartida del bucle. Esto es especialmente relevante para variables declaradas fuera del bucle y código con la semántica anterior a Go 1.22.
Ejemplo
El siguiente patrón distribuye trabajo entre un conjunto fijo de goroutines que consumen un canal compartido. Un contexto permite coordinar su cancelación. La concurrencia limitada se construye sobre esta idea.
func runWorkers(ctx context.Context, jobs <-chan Job, n int) {
var wg sync.WaitGroup
for i := 0; i < n; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for {
select {
case job, ok := <-jobs:
if !ok {
return // channel closed, no more work
}
handle(job)
case <-ctx.Done():
return // upstream cancelled
}
}
}()
}
wg.Wait()
}
Notas:
- Se crean
ntrabajadores al inicio; no se crea una goroutine nueva por cada tarea. selectsobrectx.Done()permite que los trabajadores atiendan la cancelación y terminen.sync.WaitGrouphace que quien llama espere hasta que todos los trabajadores hayan terminado. El número de goroutines que ejecutan trabajo a la vez está limitado an. La cantidad de tareas pendientes depende de la capacidad del canal y de cómo las produzca quien lo utiliza.
Propagación de errores desde goroutines
Una goroutine no devuelve valores directamente a quien la inicia. Como se ejecuta de forma independiente, sus errores deben comunicarse mediante un canal u otro mecanismo compartido. Una opción es un canal de errores con capacidad suficiente para que cada goroutine pueda enviar su error sin bloquear:
errs := make(chan error, n)
for i := 0; i < n; i++ {
go func(i int) {
if err := doWork(i); err != nil {
errs <- err
}
}(i)
}
// Collect results after all goroutines finish
// (use a WaitGroup to know when to stop reading)
wg.Wait()
close(errs)
for err := range errs {
log.Println("worker error:", err)
}
Si el canal de errores no tiene búfer y nadie recibe de él, el envío puede bloquear indefinidamente y dejar una goroutine sin terminar. Con capacidad n, cada una de las n goroutines puede enviar como máximo un error sin esperar.
Este patrón recoge los errores al terminar. Para cancelar el trabajo restante al producirse el primer error, combina el canal con la cancelación de un contexto:
ctx, cancel := context.WithCancel(context.Background())
errs := make(chan error, n)
for i := 0; i < n; i++ {
wg.Add(1)
go func(i int) {
defer wg.Done()
if err := doWork(ctx, i); err != nil {
errs <- err
cancel() // signal other goroutines to stop
}
}(i)
}
wg.Wait()
cancel() // always call cancel to release context resources
close(errs)
Hay que llamar a cancel() aunque no se produzca ningún error para liberar los recursos del contexto. Colocar defer cancel() justo después de crearlo es la forma habitual de garantizarlo.
Para flujos más complejos que deban devolver el primer error y cancelar automáticamente a los demás trabajadores, golang.org/x/sync/errgroup ofrece este patrón.
Concurrencia limitada
La concurrencia limitada establece cuántas goroutines pueden realizar trabajo simultáneamente para evitar sobrecargar la infraestructura. Aunque Go permite iniciar 50.000 goroutines con facilidad, hacerlo puede agotar las conexiones de una base de datos, superar los límites de una API o saturar la red. Hay dos patrones habituales para limitar ese trabajo:
Patrón 1: semáforo mediante un canal con búfer
Un canal con capacidad K funciona como un conjunto de permisos. Cada goroutine adquiere uno antes de trabajar y lo libera al terminar. Si los K permisos están ocupados, las siguientes deben esperar.
Es como un guardarropa con K fichas numeradas. Para trabajar hay que obtener una ficha; cuando no quedan, se espera hasta que alguien devuelva una.
sem := make(chan struct{}, K) // K = max concurrent workers
for _, job := range jobs {
sem <- struct{}{} // acquire token; blocks if K workers are already active
go func(j Job) {
defer func() { <-sem }() // release token on completion
process(j)
}(job)
}
Este patrón resulta útil para lotes breves cuando se desea limitar la concurrencia sin crear trabajadores persistentes. Se sigue creando una goroutine por tarea, pero se limita cuántas pueden realizar el trabajo a la vez.
Patrón 2: conjunto fijo de trabajadores
Se crean K goroutines al inicio. Estas esperan tareas en un canal compartido y las procesan conforme llegan. Cuando todas están ocupadas y el canal se llena, el productor debe esperar.
Es como un supermercado con K cajas: el mismo personal atiende a los clientes durante el día. Cuando todas están ocupadas, los nuevos clientes esperan; no se contrata a alguien por cada cliente.
work := make(chan Job, bufferSize)
var wg sync.WaitGroup
for i := 0; i < K; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for job := range work {
process(job)
}
}()
}
// Dispatch jobs
for _, j := range jobs {
work <- j
}
close(work) // signals workers to exit after draining
wg.Wait()
Este patrón encaja en sistemas de larga duración, como un servicio que consulta una cola continuamente. El coste de crear goroutines se paga al inicio y el tamaño del conjunto controla el trabajo simultáneo. La capacidad del búfer y el número de trabajadores son independientes. Igualarlos es una opción habitual, pero no obligatoria. Los trabajadores limitan la concurrencia; el búfer limita cuánto trabajo puede esperar en memoria. La siguiente sección explica cómo ajustar ambos. Usa un semáforo para lotes cortos con entradas conocidas y un conjunto fijo de trabajadores para servicios continuos cuyo flujo de tareas no tiene un final definido.
- Go Web Examples: Worker Pools
- Why Your Goroutines Need a Speed Limit: Bounded Concurrency in Go
- Go by Example: Worker Pools
Ajustar el tamaño del búfer
El número de trabajadores controla cuántas tareas se ejecutan a la vez. La capacidad del canal determina cuántas pueden esperar en memoria antes de bloquear al productor. No tienen por qué coincidir.
Un búfer mayor que K permite adelantar tareas mientras los trabajadores terminan las actuales. Esto reduce esperas si obtener más trabajo tiene su propia latencia, por ejemplo al consultar PostgreSQL o realizar una llamada de red.
const workers = 8
const buffer = 32 // pre-fetch up to 32 jobs ahead of workers
work := make(chan Job, buffer)
for i := 0; i < workers; i++ {
go func() {
for job := range work {
process(job)
}
}()
}
A cambio, las tareas del canal pueden haberse reclamado en el origen sin haberse completado. Si el proceso falla, su recuperación depende de cómo se reclamaron. Una transacción abierta con SELECT ... FOR UPDATE se revierte al perderse la conexión; una extracción que elimina definitivamente las tareas al leerlas puede perder las que queden en el búfer.
Un punto de partida razonable es un búfer de dos o cuatro veces el número de trabajadores, ajustándolo para mantenerlos ocupados sin retener demasiado trabajo en memoria.
Casos de uso
Extracción concurrente de páginas web
Consultar muchas URL es un caso habitual. Crear una goroutine por URL puede abrir demasiadas conexiones HTTP; además, un servidor lento puede bloquear una tarea indefinidamente si no hay un plazo de espera.
El semáforo permite limitar las peticiones simultáneas. Un context.WithTimeout por petición evita esperas indefinidas, y un canal permite recoger los resultados sin escribir simultáneamente en un segmento o mapa sin sincronización.
type Result struct {
URL string
Body []byte
Err error
}
func fetchAll(ctx context.Context, urls []string, k int) []Result {
sem := make(chan struct{}, k) // cap concurrent in-flight requests
results := make(chan Result, len(urls))
var wg sync.WaitGroup
for _, u := range urls {
wg.Add(1)
sem <- struct{}{} // acquire token; blocks if k requests are already active
go func(url string) {
defer wg.Done()
defer func() { <-sem }() // release token on completion
reqCtx, cancel := context.WithTimeout(ctx, 10*time.Second)
defer cancel()
body, err := fetch(reqCtx, url)
results <- Result{URL: url, Body: body, Err: err}
}(u)
}
wg.Wait()
close(results)
out := make([]Result, 0, len(urls))
for r := range results {
out = append(out, r)
}
return out
}
El límite k debe respetar los descriptores de archivo disponibles y los límites por servidor configurados en el transporte de http.Client. Se puede empezar entre 10 y 50 y ajustar según la tolerancia del destino y el ancho de banda.
Conjuntos de conexiones a bases de datos
Los controladores de bases de datos suelen mantener un conjunto finito de conexiones. Cada conexión consume recursos en el cliente y el servidor. Iniciar trabajo sin límites puede agotar ese conjunto y hacer que las consultas fallen o esperen. Limitar los trabajadores de acuerdo con las conexiones disponibles evita que este grupo de tareas sature el conjunto. También hay que tener en cuenta otros consumidores de esas conexiones y las tareas que necesitan más de una a la vez.
db, _ := sql.Open("postgres", connStr)
db.SetMaxOpenConns(K) // cap the pool at K connections
work := make(chan Query, bufferSize)
for i := 0; i < K; i++ {
go func() {
for q := range work {
rows, err := db.QueryContext(q.ctx, q.sql, q.args...)
// handle rows, err
}
}()
}
SetMaxOpenConns(K) y un conjunto de K trabajadores pueden coordinar ambos límites. Deben ajustarse juntos según el resto de usuarios de la base de datos y los recursos necesarios por tarea.
Límites de frecuencia y de acceso a API
Las API externas suelen limitar las peticiones por segundo, por minuto o mediante un depósito de permisos. El trabajo sin límites puede provocar errores 429 y ráfagas de reintentos que empeoran el rendimiento.
El semáforo puede combinarse con time.Ticker para controlar la frecuencia de peticiones, además de la cantidad de trabajo simultáneo.
sem := make(chan struct{}, K) // cap concurrent in-flight requests
ticker := time.NewTicker(rate) // e.g. time.Second / requestsPerSecond
defer ticker.Stop()
for _, req := range requests {
<-ticker.C // pace requests to the rate limit
sem <- struct{}{} // cap concurrent in-flight requests
go func(r Request) {
defer func() { <-sem }()
call(r)
}(req)
}
Es como un peaje con K carriles y un semáforo que regula cuándo puede entrar cada coche: el semáforo limita la frecuencia y los carriles limitan la concurrencia.
Procesamiento de archivos e imágenes
Procesar archivos, convertir imágenes o ejecutar OCR puede exigir mucha CPU y entrada/salida. Un conjunto fijo de trabajadores consume una cola y mantiene una carga predecible, sin crear miles de tareas activas a la vez.
paths := make(chan string, len(files))
for _, f := range files {
paths <- f
}
close(paths)
var wg sync.WaitGroup
for i := 0; i < K; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for path := range paths {
processFile(path)
}
}()
}
wg.Wait()
Es como un laboratorio fotográfico con K puestos de revelado. Los carretes esperan y cada puesto toma el siguiente al terminar. Añadir carretes no aumenta el número de puestos.
Para trabajo limitado por CPU, empieza con runtime.NumCPU(). Para tareas que esperan a la red o al disco, puede resultar adecuado un número mayor. Ajusta el límite mediante mediciones.
Recursos
Documentación oficial de Go
- A Tour of Go: Concurrency
- Go Dev Blog: Go Concurrency Patterns: Context
- Go Dev Blog: Share Memory by Communicating
- sync package
- context package
- errgroup package
Guías y ejemplos
- Go by Example: Channel Synchronization
- Go by Example: Worker Pools
- Go by Example: Mutexes
- Go Web Examples: Worker Pools
- Concurrency in Go: Goroutines and Channels Explained with Real Examples
- Why Your Goroutines Need a Speed Limit: Bounded Concurrency in Go
- goleak: goroutine leak detector for Go tests
Lecturas adicionales
- Concurrency in Go, de Katherine Cox-Buday: un libro sobre el modelo de concurrencia de Go que profundiza en flujos de trabajo, cancelaciones y patrones prácticos.
- El modelo de memoria de Go: define las garantías de orden que permiten coordinar el acceso concurrente de forma segura.