En los posts anteriores exploramos las goroutines, channels, select, pipelines y worker pools. Hoy vamos a profundizar en uno de los patrones más poderosos para paralelizar procesamiento: Fan-Out / Fan-In. Si vienes de Java, este patrón es similar a dividir trabajo entre múltiples threads y luego combinar resultados, pero mucho más elegante en Go.
¿Qué es Fan-Out / Fan-In?
Fan-Out y Fan-In son dos patrones complementarios que trabajan juntos:
Visualización del Patrón
Fan-Out:
Input → [Worker 1]
→ [Worker 2]
→ [Worker 3]
→ [Worker N]
Fan-In:
[Worker 1] ┐
[Worker 2] ├→ Output
[Worker 3] │
[Worker N] ┘
Comparación: Java vs Go
| Aspecto | Java | Go |
| Distribución | ExecutorService.submit() múltiples veces | Fan-out con channels |
| Combinación | Future.get() en loop | Fan-in con select/WaitGroup |
| Sincronización | Compleja (CompletableFuture, etc.) | Simple (channels) |
| Cancelación | Compleja | Simple con context |
| Expresividad | Verbosa | Simple y clara |
¿Por Qué Usar Fan-Out / Fan-In?
Problema sin paralelización:
for item := range input {
result := processExpensive(item)
output <- result
}
Problemas:
Solución con Fan-Out / Fan-In:
Ventajas:
Fan-Out: Distribuir Trabajo
Fan-Out distribuye trabajo desde un channel a múltiples workers. Cada worker procesa elementos del mismo channel.
Implementación Básica de Fan-Out
func fanOut(input <-chan int, numWorkers int) []<-chan int {
outputs := make([]<-chan int, numWorkers)
for i := 0; i < numWorkers; i++ {
output := make(chan int)
outputs[i] = output
go func(workerID int) {
defer close(output)
for item := range input {
result := processItem(item)
fmt.Printf("Worker %d processed %d -> %d\n", workerID, item, result)
output <- result
}
}(i)
}
return outputs
}
Características:
✅ Cada worker lee del mismo channel input
✅ Múltiples workers procesan en paralelo
✅ Cada worker tiene su propio channel de salida
✅ Los workers terminan cuando input se cierra
Ejemplo Completo de Fan-Out
func fanOutExample() {
input := make(chan int)
go func() {
defer close(input)
for i := 1; i <= 10; i++ {
input <- i
}
}()
numWorkers := 3
workerOutputs := make([]<-chan int, numWorkers)
for i := 0; i < numWorkers; i++ {
output := make(chan int)
workerOutputs[i] = output
go func(workerID int) {
defer close(output)
for n := range input {
time.Sleep(100 * time.Millisecond)
result := n * n
fmt.Printf("Worker %d processed %d -> %d\n", workerID, n, result)
output <- result
}
}(i)
}
}
Comportamiento:
Los workers compiten por elementos del channel input
El primer worker disponible toma el siguiente elemento
Esto distribuye el trabajo de forma balanceada
Fan-In: Combinar Resultados
Fan-In combina resultados de múltiples channels en uno solo. Hay dos formas principales de implementarlo:
Método 1: Fan-In con WaitGroup
Ideal cuando tienes un número dinámico de channels:
func fanInWaitGroup(inputs []<-chan int) <-chan int {
output := make(chan int)
var wg sync.WaitGroup
for _, input := range inputs {
wg.Add(1)
go func(ch <-chan int) {
defer wg.Done()
for item := range ch {
output <- item
}
}(input)
}
go func() {
wg.Wait()
close(output)
}()
return output
}
Ventajas:
Uso completo:
func fanOutFanInWithWaitGroup() {
input := make(chan int)
go func() {
defer close(input)
for i := 1; i <= 10; i++ {
input <- i
}
}()
numWorkers := 3
workerOutputs := make([]<-chan int, numWorkers)
for i := 0; i < numWorkers; i++ {
output := make(chan int)
workerOutputs[i] = output
go func(workerID int) {
defer close(output)
for n := range input {
time.Sleep(100 * time.Millisecond)
result := n * n
fmt.Printf("Worker %d processed %d -> %d\n", workerID, n, result)
output <- result
}
}(i)
}
results := fanInWaitGroup(workerOutputs)
for result := range results {
fmt.Printf("Final result: %d\n", result)
}
}
Método 2: Fan-In con Select
Ideal cuando tienes un número fijo y pequeño de channels:
func fanInSelect(ch1, ch2 <-chan int) <-chan int {
output := make(chan int)
go func() {
defer close(output)
ch1Active := ch1
ch2Active := ch2
for ch1Active != nil || ch2Active != nil {
select {
case val, ok := <-ch1Active:
if !ok {
ch1Active = nil
} else {
output <- val
}
case val, ok := <-ch2Active:
if !ok {
ch2Active = nil
} else {
output <- val
}
}
}
}()
return output
}
Ventajas:
✅ Más control sobre qué channel leer primero
✅ Puede priorizar channels específicos
✅ Eficiente para número pequeño de channels
Uso:
func fanOutFanInWithSelect() {
input := make(chan int)
worker1 := make(chan int)
worker2 := make(chan int)
go func() {
defer close(worker1)
for val := range input {
worker1 <- val * 2
}
}()
go func() {
defer close(worker2)
for val := range input {
worker2 <- val * 3
}
}()
output := fanInSelect(worker1, worker2)
go func() {
defer close(input)
for i := 1; i <= 5; i++ {
input <- i
}
}()
for result := range output {
fmt.Printf("Result: %d\n", result)
}
}
Comparación: WaitGroup vs Select
| Aspecto | WaitGroup | Select |
| Número de channels | Dinámico (cualquier número) | Fijo (pequeño número) |
| Priorización | No | Sí (orden de cases) |
| Complejidad | Baja | Media |
| Uso de memoria | Más goroutines | Menos goroutines |
| Cuándo usar | Número variable de workers | 2-5 channels fijos |
Ejemplo Práctico: Procesar Archivos en Paralelo
Un caso de uso real: procesar múltiples archivos en paralelo y combinar resultados:
type FileResult struct {
Filename string
Lines int
Error error
}
func processFilesInParallel(filenames []string, numWorkers int) <-chan FileResult {
files := make(chan string)
go func() {
defer close(files)
for _, filename := range filenames {
files <- filename
}
}()
workerOutputs := make([]<-chan FileResult, numWorkers)
for i := 0; i < numWorkers; i++ {
output := make(chan FileResult)
workerOutputs[i] = output
go func(workerID int) {
defer close(output)
for filename := range files {
lines, err := countLines(filename)
output <- FileResult{
Filename: filename,
Lines: lines,
Error: err,
}
}
}(i)
}
results := make(chan FileResult)
var wg sync.WaitGroup
for _, workerOutput := range workerOutputs {
wg.Add(1)
go func(ch <-chan FileResult) {
defer wg.Done()
for result := range ch {
results <- result
}
}(workerOutput)
}
go func() {
wg.Wait()
close(results)
}()
return results
}
Fan-Out / Fan-In con Context (Cancelación)
Agregar cancelación hace el patrón más robusto:
func fanOutFanInWithContext(ctx context.Context, input <-chan int, numWorkers int) <-chan int {
workerOutputs := make([]<-chan int, numWorkers)
for i := 0; i < numWorkers; i++ {
output := make(chan int)
workerOutputs[i] = output
go func(workerID int) {
defer close(output)
for {
select {
case <-ctx.Done():
return
case item, ok := <-input:
if !ok {
return
}
result := processItem(ctx, item)
select {
case output <- result:
case <-ctx.Done():
return
}
}
}
}(i)
}
results := make(chan int)
var wg sync.WaitGroup
for _, workerOutput := range workerOutputs {
wg.Add(1)
go func(ch <-chan int) {
defer wg.Done()
for {
select {
case <-ctx.Done():
return
case result, ok := <-ch:
if !ok {
return
}
select {
case results <- result:
case <-ctx.Done():
return
}
}
}
}(workerOutput)
}
go func() {
wg.Wait()
close(results)
}()
return results
}
Fan-Out con Rate Limiting
Limitar el rate de distribución puede ser útil:
func fanOutWithRateLimit(input <-chan int, numWorkers int, rate time.Duration) []<-chan int {
outputs := make([]<-chan int, numWorkers)
limiter := make(chan struct{}, numWorkers)
go func() {
ticker := time.NewTicker(rate)
defer ticker.Stop()
for range ticker.C {
select {
case limiter <- struct{}{}:
default:
}
}
}()
for i := 0; i < numWorkers; i++ {
output := make(chan int)
outputs[i] = output
go func(workerID int) {
defer close(output)
for item := range input {
<-limiter
result := processItem(item)
output <- result
}
}(i)
}
return outputs
}
Fan-In con Prioridades
Cuando necesitas priorizar ciertos channels:
func fanInWithPriority(highPriority, lowPriority <-chan int) <-chan int {
output := make(chan int)
go func() {
defer close(output)
highActive := highPriority
lowActive := lowPriority
for highActive != nil || lowActive != nil {
select {
case val, ok := <-highActive:
if !ok {
highActive = nil
} else {
output <- val
}
case val, ok := <-lowActive:
if !ok {
lowActive = nil
} else {
output <- val
}
}
}
}()
return output
}
Ejemplo Avanzado: Pipeline con Fan-Out/Fan-In
Combinando pipelines con fan-out/fan-in:
func advancedPipeline() {
input := make(chan int)
go func() {
defer close(input)
for i := 1; i <= 100; i++ {
input <- i
}
}()
numTransformers := 5
transformedOutputs := make([]<-chan int, numTransformers)
for i := 0; i < numTransformers; i++ {
output := make(chan int)
transformedOutputs[i] = output
go func() {
defer close(output)
for n := range input {
output <- transform(n)
}
}()
}
transformed := fanInWaitGroup(transformedOutputs)
filtered := make(chan int)
go func() {
defer close(filtered)
for n := range transformed {
if n > 50 {
filtered <- n
}
}
}()
for result := range filtered {
fmt.Printf("Result: %d\n", result)
}
}
Mejores Prácticas
1. Elegir el Número Correcto de Workers
numWorkers := runtime.NumCPU()
numWorkers := runtime.NumCPU() * 2
numWorkers := getConfig().WorkerCount
numWorkers := 1000
2. Usar WaitGroup para Fan-In Dinámico
func fanInWaitGroup(inputs []<-chan int) <-chan int {
output := make(chan int)
var wg sync.WaitGroup
for _, input := range inputs {
wg.Add(1)
go func(ch <-chan int) {
defer wg.Done()
for item := range ch {
output <- item
}
}(input)
}
go func() {
wg.Wait()
close(output)
}()
return output
}
3. Usar Select para Fan-In de Pocos Channels
func fanInSelect(ch1, ch2 <-chan int) <-chan int {
}
func fanInSelectMany(ch1, ch2, ch3, ch4, ch5, ch6, ch7, ch8 <-chan int) <-chan int {
}
4. Siempre Cerrar Channels Correctamente
go func() {
defer close(output)
for item := range input {
output <- process(item)
}
}()
go func() {
for item := range input {
output <- process(item)
}
}()
5. Manejar Errores en Fan-Out
type Result struct {
Value int
Error error
}
func fanOutWithErrors(input <-chan int) []<-chan Result {
outputs := make([]<-chan Result, numWorkers)
for i := 0; i < numWorkers; i++ {
output := make(chan Result)
outputs[i] = output
go func() {
defer close(output)
for item := range input {
result, err := processWithError(item)
output <- Result{Value: result, Error: err}
}
}()
}
return outputs
}
6. Usar Buffers Apropiados
output := make(chan int, numWorkers*2)
output := make(chan int)
output := make(chan int, 1000000)
Errores Comunes
❌ Error 1: Leak de Goroutines en Fan-In
func badFanIn(inputs []<-chan int) <-chan int {
output := make(chan int)
for _, input := range inputs {
go func(ch <-chan int) {
for item := range ch {
output <- item
}
}(input)
}
return output
}
func goodFanIn(inputs []<-chan int) <-chan int {
output := make(chan int)
var wg sync.WaitGroup
for _, input := range inputs {
wg.Add(1)
go func(ch <-chan int) {
defer wg.Done()
for item := range ch {
output <- item
}
}(input)
}
go func() {
wg.Wait()
close(output)
}()
return output
}
❌ Error 2: No Verificar Cierre en Select
select {
case val := <-ch1:
process(val)
case val := <-ch2:
process(val)
}
ch1Active := ch1
ch2Active := ch2
for ch1Active != nil || ch2Active != nil {
select {
case val, ok := <-ch1Active:
if !ok {
ch1Active = nil
} else {
process(val)
}
case val, ok := <-ch2Active:
if !ok {
ch2Active = nil
} else {
process(val)
}
}
}
❌ Error 3: Race Condition en Fan-Out
for i := 0; i < numWorkers; i++ {
go func() {
fmt.Printf("Worker %d\n", i)
}()
}
for i := 0; i < numWorkers; i++ {
go func(workerID int) {
fmt.Printf("Worker %d\n", workerID)
}(i)
}
❌ Error 4: Demasiados o Muy Pocos Workers
numWorkers := 10000
numWorkers := 1
numWorkers := runtime.NumCPU()
❌ Error 5: No Manejar Backpressure
for item := range input {
output <- process(item)
}
for item := range input {
result := process(item)
select {
case output <- result:
case <-ctx.Done():
return
default:
log.Printf("Output full, dropping result")
}
}
Comparación: Java vs Go
Java: ExecutorService y CompletableFuture
ExecutorService executor = Executors.newFixedThreadPool(10);
List<CompletableFuture<String>> futures = new ArrayList<>();
for (String file : files) {
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
return processFile(file);
}, executor);
futures.add(future);
}
CompletableFuture<Void> allFutures = CompletableFuture.allOf(
futures.toArray(new CompletableFuture[0])
);
List<String> results = allFutures.thenApply(v -> {
return futures.stream()
.map(CompletableFuture::join)
.collect(Collectors.toList());
}).join();
Problemas:
Go: Fan-Out / Fan-In
input := make(chan string)
go func() {
defer close(input)
for _, file := range files {
input <- file
}
}()
numWorkers := 10
workerOutputs := make([]<-chan string, numWorkers)
for i := 0; i < numWorkers; i++ {
output := make(chan string)
workerOutputs[i] = output
go func() {
defer close(output)
for file := range input {
output <- processFile(file)
}
}()
}
results := fanInWaitGroup(workerOutputs)
for result := range results {
fmt.Println(result)
}
Ventajas:
Conclusiones
Fan-Out / Fan-In es un patrón esencial en Go:
✅ Paralelización: Distribuye trabajo eficientemente a múltiples workers
✅ Agregación: Combina resultados de forma elegante
✅ Escalabilidad: Mejora el throughput significativamente
✅ Flexibilidad: Dos métodos (WaitGroup y Select) para diferentes necesidades
✅ Cancelación: Fácil agregar cancelación con context
✅ Simplicidad: Código claro y mantenible
Si vienes de Java, verás que Fan-Out / Fan-In en Go es mucho más simple y expresivo que usar ExecutorService y CompletableFuture. La combinación de goroutines y channels hace que paralelizar trabajo sea natural y seguro.
Próximos Pasos
En el siguiente post exploraremos otros patrones avanzados de concurrencia en Go, incluyendo circuit breakers, retry/backoff, y otros patrones que complementan los que hemos aprendido.
¿Has usado Fan-Out / Fan-In en tus proyectos de Go? ¿Prefieres WaitGroup o Select? Comparte tus experiencias y casos de uso en los comentarios. Y si quieres ver el código completo de estos ejemplos, puedes encontrarlo en mi repositorio go-mastery-lab.