🔀 Fan-Out / Fan-In: Строим конвейер, который не лопнет Представьте задачу: у вас есть CSV-файл на 10 миллионов строк (и
Представьте задачу: у вас есть CSV-файл на 10 миллионов строк (или бесконечный стрим из Kafka). Каждую строку нужно прочитать, сходить с ней в тяжелый внешний API (парсинг/обогащение) и записать результат в базу.
• Решение джуна: Читать по одной строке, ходить в API, писать в БД. Очень надежно и очень медленно. Файл будет обрабатываться неделю.
• Решение мидла: На каждую строку делать go func(). Через секунду мы откроем 10 миллионов горутин, забьем сеть, положим внешний API, исчерпаем файловые дескрипторы и умрем от OOM (Out Of Memory).
Нам нужен баланс: обрабатывать данные параллельно, но с жестким лимитом ресурсов. Встречайте паттерн Pipeline (Конвейер) с применением Fan-Out / Fan-In.
Что это такое?
• Fan-Out (Разветвление): Один канал генерирует задачи, а группа из N воркеров (фиксированный пул) читает из этого одного канала. Задачи распределяются между ними автоматически.
• Fan-In (Слияние): Воркеры пишут результаты в свои личные исходящие каналы, а специальная функция сливает эти N каналов в один итоговый поток.
Как это выглядит в коде (The Go Way):
// 1. Fan-Out: Воркер читает из in и пишет в свой out
func worker(in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for n := range in {
// Имитируем тяжелую работу
time.Sleep(time.Millisecond * 100)
out <- n * 2
}
}()
return out
}
// 2. Fan-In: Сливаем каналы от всех воркеров в один
func merge(cs ...<-chan int) <-chan int {
var wg sync.WaitGroup
out := make(chan int)
// Функция, которая перекладывает данные из конкретного канала в общий
output := func(c <-chan int) {
defer wg.Done()
for n := range c {
out <- n
}
}
wg.Add(len(cs))
for _, c := range cs {
go output(c)
}
// Фоновая горутина закроет общий канал, когда все воркеры отработают
go func() {
wg.Wait()
close(out)
}()
return out
}
Собираем всё вместе:
func main() {
// Канал с задачами (генератор опустим для краткости)
in := generateTasks()
// Запускаем Fan-Out: создаем фиксированно 3 воркера
w1 := worker(in)
w2 := worker(in)
w3 := worker(in)
// Запускаем Fan-In: собираем результаты из 3 каналов в 1
for result := range merge(w1, w2, w3) {
fmt.Println(result)
}
}
🔥 Нюансы для Senior-ов:
1. Почему просто не писать всем воркерам в один общий канал?
Можно. Часто так и делают (называется Worker Pool). Но классический Fan-In (с функцией merge) дает гибкость: вы можете строить сложные графы обработки, где каналы передаются из функции в функцию, не завязываясь на глобальные состояния и мьютексы.
2. Утечки горутин (Goroutine Leaks).
В этом коде есть слабое место. Если цикл чтения итогового результата в main прервется досрочно (например, возникла ошибка и мы сделали return или break), воркеры зависнут навсегда, пытаясь записать данные в каналы, которые никто не читает.
Золотое правило: Всегда прокидывайте context.Context или канал done во все функции конвейера и проверяйте case <-ctx.Done(): внутри циклов for.
3. Порядок не гарантирован.
Fan-Out перемешивает данные. Если вам критически важно сохранить исходную последовательность строк CSV, этот паттерн нужно усложнять (например, передавать структуру с индексом и сортировать буфер на выходе).
#golang #concurrency #architecture #patterns #cleancode
👉 @golang_lib