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.

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 n trabajadores al inicio; no se crea una goroutine nueva por cada tarea.
  • select sobre ctx.Done() permite que los trabajadores atiendan la cancelación y terminen.
  • sync.WaitGroup hace que quien llama espere hasta que todos los trabajadores hayan terminado. El número de goroutines que ejecutan trabajo a la vez está limitado a n. 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.

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

Guías y ejemplos

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.