En el desarrollo de sistemas de alto rendimiento, un escenario habitual es la necesidad de consolidar múltiples fuentes de datos para responder a una petición. Imaginemos el ingreso a una sala de transmisión en vivo: se requiere recuperar datos del emisor, estadísticas de audiencia, configuraciones de la sala, privilegios del espectador e información de la plataforma.
Bajo un enfoque tradicional sin concurrencia, estas consultas se ejecutan una tras otra. El tiempo total de respuesta equivale a la suma de todas las latencias individuales. Para cinco operaciones que tardan 1, 2, 3, 4 y 5 segundos respectivamente, el usuario esperaría 15 segundos, lo cual es inaceptable.
Patrón 1: Sincronización con sync.WaitGroup
Go proporciona sync.WaitGroup como mecanismo primitivo para coordinar goroutines. Permite bloquear la ejecución hasta que un conjunto de tareas concurrentes finalice.
package main
import (
"context"
"fmt"
"sync"
"time"
)
func main() {
var (
platformInfo int
streamerData int
audienceStats int
roomConfig int
privileges int
)
ctx := context.Background()
ExecuteParallel(ctx,
func() {
platformInfo = fetchPlatformData()
},
func() {
streamerData = fetchStreamerProfile()
},
func() {
audienceStats = calculateAudienceMetrics()
},
func() {
roomConfig = loadRoomSettings()
},
func() {
privileges = resolveUserPrivileges()
},
)
fmt.Println(platformInfo, streamerData, audienceStats, roomConfig, privileges)
}
func ExecuteParallel(ctx context.Context, tasks ...func()) {
var barrier sync.WaitGroup
for _, task := range tasks {
barrier.Add(1)
go func(t func()) {
defer barrier.Done()
t()
}(task)
}
barrier.Wait()
}
func fetchPlatformData() int {
time.Sleep(1 * time.Second)
fmt.Println("Consulta de plataforma completada")
return 100
}
func fetchStreamerProfile() int {
time.Sleep(2 * time.Second)
fmt.Println("Perfil del emisor recuperado")
return 200
}
func calculateAudienceMetrics() int {
time.Sleep(3 * time.Second)
fmt.Println("Métricas de audiencia calculadas")
return 300
}
func loadRoomSettings() int {
time.Sleep(4 * time.Second)
fmt.Println("Configuración de sala cargada")
return 400
}
func resolveUserPrivileges() int {
time.Sleep(5 * time.Second)
fmt.Println("Privilegios de usuario resueltos")
return 500
}
Con este patrón, la duración total se reduce a la tarea más lenta (5 segundos), independientemente de cuántas operaciones se ejecuten en paralelo.
Patrón 2: Manejo de errores con errgroup
El paquete golang.org/x/sync/errgroup extiende la funcionalidad de WaitGroup incorporando propagación de errores y cancelación contextual. Es particularmente útil cuando las tareas pueden fallar y se necesita detectar la primera anomalía.
package main
import (
"context"
"errors"
"fmt"
"time"
"golang.org/x/sync/errgroup"
)
func main() {
var (
platformInfo int
streamerData int
audienceStats int
roomConfig int
privileges int
)
ctx := context.Background()
err := ExecuteWithErrorHandling(ctx,
func() error {
platformInfo = fetchPlatformData()
return nil
},
func() error {
streamerData = fetchStreamerProfile()
return nil
},
func() error {
audienceStats = calculateAudienceMetrics()
return errors.New("fallo en servicio de métricas")
},
func() error {
roomConfig = loadRoomSettings()
return nil
},
func() error {
privileges = resolveUserPrivileges()
return nil
},
)
if err != nil {
fmt.Printf("Error durante ejecución concurrente: %v\n", err)
return
}
fmt.Println(platformInfo, streamerData, audienceStats, roomConfig, privileges)
}
func ExecuteWithErrorHandling(ctx context.Context, operations ...func() error) error {
var group errgroup.Group
for idx := range operations {
operation := operations[idx]
group.Go(func() error {
return operation()
})
}
return group.Wait()
}
La trampa del cierre sobre variables de ciclo
Un error recurrente al usar errgroup es capturar incorrectamente la variable del iterador. Consideremos tres variantes:
Variante defectuosa: captura directa
for _, op := range operations {
group.Go(func() error {
return op() // ¡Todas las goroutines ven la misma 'op'!
})
}
Debido a que la variable op se reutiliza en cada iteración, todas las goroutines acceden al último valor asignado, provocando comportamiento errático o panics.
Variante correcta A: reasignación explícita
for i := range operations {
op := operations[i]
group.Go(func() error {
return op()
})
}
Variante correcta B: copia en ámbito local
for _, op := range operations {
currentOp := op
group.Go(func() error {
return currentOp()
})
}
Ambas soluciones garantizan que cada clausura capture una instancia distinta de la función, eliminando la condición de carrera.
Consideraciones de diseño
Al implementar patrones concurrentes, es fundamental evaluar:
- Granularidad: Dividir excesivamente aumenta la sobrecarga de sincronización. Agrupar operaciones dependientes mejora la eficiencia.
- Cancelación: Propagar
context.Contextpermite interrumpir goroutines cuando el cliente abandona la petición. - Límites de concurrencia: Para proteger servicios externos, considerar
errgroup.Groupcon semáforos ogolang.org/x/sync/semaphore. - Idempotencia: Las tareas deben tolerar reintentos sin efectos secundarios indeseados.
La elección entre sync.WaitGroup y errgroup depende de si se requiere detección temprana de fallos. Para pipelines donde cualquier error invalida el resultado completo, errgroup es preferible. Para recopilación de métricas donde se desea completitud parcial, WaitGroup con canales de error manuales ofrece mayor flexibilidad.