Lanzar una goroutine por cada tarea es fácil y casi siempre mala idea: con diez mil tareas tendrás diez mil peticiones simultáneas. Lo que se quiere es un número fijo de trabajadores que van tomando tareas de una cola. En Go eso son unas veinte líneas.
También te puede interesar
El punto de partida
Nueve tareas que tardan 100 ms cada una —descargar algo, convertir una imagen, consultar una API—. Una detrás de otra:
func procesar(t Tarea) (int, error) {
time.Sleep(100 * time.Millisecond)
if t.ID == 4 {
return 0, fmt.Errorf("tarea %d: el servidor respondió 500", t.ID)
}
return t.ID * 1000, nil
}
inicio := time.Now()
for _, t := range tareas {
procesar(t)
}
fmt.Printf("secuencial: %v\n", time.Since(inicio).Round(10*time.Millisecond))
salidasecuencial: 900ms
El grupo de trabajadores
Tres piezas: una cola de entrada, un canal de salida y un WaitGroup que sabe cuándo han terminado todos.
const nObreros = 3
cola := make(chan Tarea)
salida := make(chan Resultado)
var wg sync.WaitGroup
for o := 1; o <= nObreros; o++ {
wg.Add(1)
go func(obrero int) {
defer wg.Done()
for t := range cola { // toma tareas hasta que se cierre la cola
n, err := procesar(t)
salida <- Resultado{ID: t.ID, Bytes: n, Err: err, Obrero: obrero}
}
}(o)
}
// el que reparte
go func() {
for _, t := range tareas {
cola <- t
}
close(cola) // así los trabajadores salen del for range
}()
// el que cierra la salida cuando ya no queda nadie trabajando
go func() {
wg.Wait()
close(salida)
}()
for r := range salida {
// … contar, sumar, guardar
}
salidasecuencial: 900ms fallo: tarea 4: el servidor respondió 500 con 3 trabajadores: 300ms 8 bien, 1 con error, 41000 bytes en total trabajador 1 atendió 3 tareas trabajador 2 atendió 3 tareas trabajador 3 atendió 3 tareas
900 ms a 300: exactamente el triple, porque son tres trabajadores y las tareas duran todas lo mismo. El reparto salió 3-3-3 en esta ejecución; con tareas de duración desigual no sería tan redondo, y eso es precisamente lo bueno del patrón —el que acaba antes coge la siguiente, sin que nadie tenga que calcular el reparto.
cola no tiene búfer: el primer envío se quedaría bloqueado esperando a un trabajador que todavía no existe. Y el cierre de salida va en otra porque wg.Wait() bloquea hasta que todos acaben, y nadie acabará mientras el bucle principal no lea de salida. Si pones cualquiera de las dos en línea, el programa se cuelga: fatal error: all goroutines are asleep - deadlock!Cuántos trabajadores
Depende de a qué esperan. Si el cuello de botella es la red o el disco, pueden ser decenas o cientos, porque pasan el rato dormidas. Si es cálculo puro, el número útil es runtime.NumCPU(): más no reparte mejor, solo añade cambios de contexto.
Por qué no se suma a una variable compartida
La tentación es prescindir del canal de salida y que cada trabajador vaya sumando en una variable:
total := 0
for i := 1; i <= 1000; i++ {
wg.Add(1)
go func() {
defer wg.Done()
total++ // ¡mal! dos goroutines pueden hacerlo a la vez
}()
}
wg.Wait()
fmt.Println("total:", total, "(deberían ser 1000)")
Tres ejecuciones seguidas del mismo programa:
salidatotal: 994 (deberían ser 1000) total: 988 (deberían ser 1000) total: 971 (deberían ser 1000)
Nunca da 1000 y nunca da lo mismo. total++ no es una sola operación: lee, suma y escribe. Si dos goroutines leen el mismo valor antes de que ninguna escriba, uno de los dos incrementos se pierde.
Lo peor no es el resultado: es que el programa no se queja. Por eso existe el detector de carreras, que se activa con una bandera:
go run -race .
salida con -race==================
WARNING: DATA RACE
Read at 0x00c00011a038 by goroutine 9:
main.main.func1()
/home/ana/w/carrera/main.go:15 +0x7b
Previous write at 0x00c00011a038 by goroutine 8:
main.main.func1()
/home/ana/w/carrera/main.go:15 +0x8d
Goroutine 9 (running) created at:
main.main()
/home/ana/w/carrera/main.go:13 +0x8a
Te da la línea exacta y las dos goroutines implicadas. Acostúmbrate a pasar -race a las pruebas (go test -race ./...): encuentra en segundos cosas que de otro modo aparecen una vez al mes en producción. No se usa en la versión final porque hace el programa bastante más lento y pesado.
¿La solución? Un sync.Mutex alrededor del incremento, o un atomic.Int64, o —lo más idiomático en Go— no compartir la variable y mandar los resultados por un canal, como en el grupo de trabajadores de arriba.
Parar en cuanto algo falla
Con veinte tareas y un error en la cuarta, no tiene sentido seguir con las dieciséis restantes. Eso lo resuelve un context:
func procesar(ctx context.Context, id int) (int, error) {
select {
case <-time.After(100 * time.Millisecond): // el trabajo de verdad
case <-ctx.Done():
return 0, ctx.Err() // alguien canceló
}
...
}
El trabajador se retira al primer error, y el repartidor también mira el contexto para no quedarse bloqueado enviando a una cola que ya no lee nadie:
// repartidor
for id := 1; id <= 20; id++ {
select {
case cola <- id:
case <-ctx.Done():
fmt.Println("el repartidor para en la tarea", id)
return
}
}
var primero error
for err := range salida {
if errors.Is(err, context.Canceled) {
continue // eco de la cancelación, no un fallo de verdad
}
if primero == nil {
primero = err
cancelar() // avisa a todos los demás
}
}
salidael repartidor para en la tarea 9 error: tarea 4: respuesta 500 del servidor terminado en 200ms
Veinte tareas a 100 ms entre tres trabajadores habrían sido 700 ms. Paró en 200. Fíjate en el filtro de context.Canceled: al cancelar, las tareas en vuelo devuelven ese error, y si no lo descartas acabas con una lista de «errores» que en realidad son el eco de tu propia cancelación.
defer cancelar() justo después de crear el contexto no es opcional aunque ya llames a cancelar() en el bucle: si el programa sale por otro camino, el contexto se queda sin liberar. go vet avisa de ello.Las reglas que evitan los bloqueos
| Regla | Por qué |
|---|---|
| Cierra el canal quien escribe en él | Escribir en un canal cerrado es panic |
| Cierra la cola cuando no haya más tareas | Si no, los for range no terminan nunca |
wg.Add antes de lanzar la goroutine |
Dentro, puede ejecutarse después del Wait |
defer wg.Done(), siempre |
Un return temprano dejaría el Wait colgado |
El canal de salida se cierra tras wg.Wait() |
Y ese Wait va en su propia goroutine |
Prueba con go test -race |
Es lo único que encuentra las carreras |
Todo el código se compiló y ejecutó con Go 1.25 antes de publicarla; los tiempos y la salida del detector de carreras son los reales de esa ejecución. Documentación oficial: sync.WaitGroup y el detector de carreras.