Отдельные примитивы — горутины, каналы,
select, context — это кирпичи. Сами по себе они
мало что значат: интересное начинается, когда из них складывают устойчивые
конструкции, повторяющиеся в любом проде. Пайплайн, который тянет данные через
несколько стадий. Пул воркеров, который держит нагрузку в узде. Семафор, который
не даёт открыть тысячу соединений разом. Батчер, который копит мелочь и пишет в
базу пачками.
Эта глава — каталог таких паттернов. Но не «вот код, копируй». Главная мысль
сквозная и простая: любой конкурентный граф корректен ровно настолько,
насколько чётко в нём отвечены два вопроса — как до каждого узла доходит сигнал
«данные кончились / пора останавливаться» и кто закрывает каждый канал. Если на
оба есть ответ — структура не утечёт. Если хоть на один
ответа нет — рано или поздно повиснет горутина.
Примеры ниже по возможности запускаемые прямо в браузере. Там, где речь о
настоящем параллелизме или гонках, — статичный код с разбором, потому что в
браузере всё исполняется однопоточно и race-детектора нет.
Пайплайн — это конвейер. Цепочка стадий, где выход одной становится входом
следующей. Каждая стадия — горутина: читает из входного канала, что-то делает,
пишет в выходной. Стадии работают одновременно, и пока последняя печатает первый
результат, первая уже готовит десятый.
Ключевая дисциплина одна: стадия закрывает свой выходной канал, когда её вход
иссяк. Обычно через defer close(out). Это и есть способ передать «данные
кончились» вниз по конвейеру — следующая стадия увидит закрытие через range и
сама закроется. Сигнал конца течёт по тем же рельсам, что и данные.
package mainimport "fmt"// gen — источник: выкладывает числа и закрывает выходfunc gen(nums ...int) <-chan int { out := make(chan int) go func() { defer close(out) for _, n := range nums { out <- n } }() return out}// sq — стадия: возводит в квадрат всё, что пришлоfunc sq(in <-chan int) <-chan int { out := make(chan int) go func() { defer close(out) for n := range in { out <- n * n } }() return out}func main() { // gen -> sq -> потребитель for v := range sq(gen(1, 2, 3, 4)) { fmt.Println(v) }}
Обрати внимание: ни одна стадия не знает, сколько чисел придёт и когда. Она просто
читает, пока вход открыт, и закрывает свой выход, когда range завершился. Можно
добавить третью стадию, четвёртую — логика каждой не меняется.
Пока потребитель честно вычитывает всё до конца — проблем нет. Беда приходит,
когда он уходит раньше: взял первый результат и сделал break. Стадия выше
застрянет на out <- v, потому что читать её выход больше некому. Горутина повисла
навсегда — классическая утечка.
Лекарство — дать стадии способ услышать «всё, расходимся». Это context:
на каждой отправке слушаем ещё и ctx.Done().
func sq(ctx context.Context, in <-chan int) <-chan int { out := make(chan int) go func() { defer close(out) for n := range in { select { case out <- n * n: case <-ctx.Done(): return // потребитель ушёл — не виснем на отправке } } }() return out}
Без ветки <-ctx.Done() отправка out <- v зависла бы навсегда, если ниже по
конвейеру перестали читать. С ней — горутина аккуратно выходит, defer close(out)
отрабатывает, сигнал отмены катится дальше вниз. Запомни это как рефлекс: каждая
отправка в долгоживущий канал получает соседнюю ветку с отменой.
Иногда одна стадия — узкое место: данные простые, а их много. Тогда её
распараллеливают.
Fan-out — несколько горутин читают из одного входного канала. Каждая берёт
следующий элемент, как только освободилась; быстрые забирают больше, медленные
меньше — нагрузка балансируется сама. Fan-in — обратная операция: слить
несколько выходных каналов обратно в один, чтобы потребитель видел единый поток.
Fan-in — это та самая задача «несколько отправителей, один канал» из главы про
каналы. И тонкость там ровно одна: выходной канал закрывает
кто-то один, и только когда замолчали все источники. Иначе либо запись в
закрытый канал (паника), либо double-close (тоже паника). Идиома — sync.WaitGroup:
каждый источник делает Done, отдельная горутина ждёт Wait и закрывает выход.
package mainimport ( "fmt" "sort" "sync")func gen(nums ...int) <-chan int { out := make(chan int) go func() { defer close(out) for _, n := range nums { out <- n } }() return out}// один воркер fan-out: квадраты из общего входаfunc worker(in <-chan int) <-chan int { out := make(chan int) go func() { defer close(out) for n := range in { out <- n * n } }() return out}// fan-in: слить каналы в один, закрыть после Waitfunc merge(chans ...<-chan int) <-chan int { out := make(chan int) var wg sync.WaitGroup wg.Add(len(chans)) for i := 0; i < len(chans); i++ { c := chans[i] go func() { defer wg.Done() for v := range c { out <- v } }() } go func() { wg.Wait() // дождались всех источников close(out) // и только теперь закрыли выход — закрывает ОДИН владелец }() return out}func main() { in := gen(1, 2, 3, 4, 5, 6) // fan-out: три воркера разбирают один вход w1, w2, w3 := worker(in), worker(in), worker(in) // fan-in: собрали обратно var got []int for v := range merge(w1, w2, w3) { got = append(got, v) } // порядок недетерминирован — сортируем, чтобы вывод был стабильным sort.Ints(got) fmt.Println(got)}
Здесь видно, почему результат приходится сортировать перед печатью: три воркера
работают независимо, и кто допишет первым — заранее не известно. Fan-out почти
всегда перемешивает порядок. Если порядок важен — либо тащите его через сам
элемент (порядковый номер), либо не распараллеливайте эту стадию.
Боевой merge так же добавляет select с <-ctx.Done() на отправке в out —
по той же причине, что и в пайплайне: чтобы источники не зависли, если потребитель
ушёл раньше.
Fan-out с каналом квадратов — это уже почти пул. Worker pool делает идею явной:
фиксированное число воркеров разбирает задачи из общей очереди и складывает
результаты в общий канал. Зачем фиксированное? Потому что «горутина на задачу»
при миллионе задач — это миллион одновременных обращений к базе или диску. Пул из
восьми воркеров обрабатывает всё то же, но держит конкуренцию ровно на восьми.
Размер пула — ваш главный рычаг backpressure.
package mainimport ( "context" "fmt" "sort" "sync")func pool(ctx context.Context, jobs <-chan int, n int) <-chan int { results := make(chan int) var wg sync.WaitGroup for i := 0; i < n; i++ { wg.Add(1) go func() { defer wg.Done() for { select { case <-ctx.Done(): return case j, ok := <-jobs: if !ok { return // очередь закрыта — воркер выходит } select { case results <- j * j: case <-ctx.Done(): return } } } }() } go func() { wg.Wait() close(results) // закрывает один владелец после всех воркеров }() return results}func main() { ctx := context.Background() jobs := make(chan int) go func() { defer close(jobs) // поставщик закрывает очередь — это сигнал воркерам for i := 1; i <= 6; i++ { jobs <- i } }() var got []int for r := range pool(ctx, jobs, 3) { // 3 воркера на 6 задач got = append(got, r) } sort.Ints(got) fmt.Println(got)}
Разберём завершение, потому что в нём вся соль. Поставщик закрывает jobs —
каждый воркер ловит ok == false и выходит. Когда вышли все, wg.Wait()
разблокируется и закрывает results. Потребительский range видит закрытие и
останавливается. Ни одна горутина не повисла: у каждой есть и причина выйти
(jobs закрыт), и аварийный выход (ctx.Done()).
Заметь симметрию со всеми примерами выше — wg.Wait() закрывает агрегирующий
канал, а не сами воркеры закрывают results. Воркеров много, владелец закрытия —
один.
Пул хорош, когда есть поток однотипных задач и канал результатов. Но иногда нужно
проще: «запусти вот эти задачи, но не больше N одновременно». Заводить ради этого
пул и канал результатов — перебор. Хватает буферизованного канала как счётного
семафора.
Идея игрушечная: канал на N элементов — это N пропусков. Чтобы стартовать,
горутина кладёт пропуск в канал (заняла слот). Если все N слотов заняты —
отправка блокируется, и новая горутина ждёт. Закончила — забирает пропуск обратно
(освободила слот).
package mainimport ( "fmt" "sync/atomic")func main() { const maxConcurrency = 2 sem := make(chan struct{}, maxConcurrency) var inFlight int32 // сколько слотов занято прямо сейчас for i := 1; i <= 6; i++ { sem <- struct{}{} // занять слот (блок, если все заняты) held := atomic.AddInt32(&inFlight, 1) fmt.Printf("задача %d: занято слотов %d/%d\n", i, held, maxConcurrency) atomic.AddInt32(&inFlight, -1) <-sem // освободить слот } fmt.Println("все 6 задач прошли через", maxConcurrency, "слота")}
Семафор пропускает задачи порциями по cap(sem): пока слот занят, попытка занять
сверх лимита (sem <- struct{}{}) заблокировалась бы. Полный буфер — естественный
backpressure: новые горутины просто ждут на отправке в канал, пока кто-то не
освободит слот через <-sem. Никаких счётчиков и мьютексов руками.
В проде вместо самодельного канала обычно берут golang.org/x/sync/semaphore —
там взвешенные слоты (одна задача может «весить» больше одной) и Acquire с
поддержкой context, чтобы ожидание слота было отменяемым.
Последний паттерн — про экономию. Писать в базу по одной записи дорого: каждая —
round-trip. Батчинг копит элементы и сбрасывает их пачкой. Но если просто ждать,
пока наберётся size, можно зависнуть: пришло три элемента, поток затих — и они
никогда не запишутся. Поэтому у батчера два триггера: набралось size или
вышло время, что наступит раньше. Оба живут в одном select.
package mainimport ( "fmt" "time")func main() { in := make(chan int) done := make(chan struct{}) go func() { defer close(done) const size = 3 buf := make([]int, 0, size) timer := time.NewTimer(50 * time.Millisecond) defer timer.Stop() flush := func(reason string) { if len(buf) > 0 { fmt.Printf("flush(%s): %v\n", reason, buf) buf = buf[:0] } timer.Reset(50 * time.Millisecond) } for { select { case it, ok := <-in: if !ok { flush("close") return } buf = append(buf, it) if len(buf) >= size { flush("size") // набрали пачку } case <-timer.C: flush("time") // не набрали, но вышло время } } }() // 4 элемента быстро: первая пачка по size, остаток уйдёт по close for i := 1; i <= 4; i++ { in <- i } close(in) <-done}
Здесь первая пачка [1 2 3] уходит по триггеру size, а хвост [4] — по
закрытию входа. В реальном сервисе срабатывал бы ещё и таймер: поток замолчал, а
накопленное всё равно надо сбросить, не дожидаясь следующего элемента. Боевой
батчер также слушает ctx.Done(), чтобы при остановке сервиса сделать финальный
flush и не потерять накопленное.
Тонкость с таймером: после каждого flush его надо Reset, иначе окно времени
поедет. И не забывайте defer timer.Stop() — иначе утечёт сам таймер.
Очень частая задача: запустить N независимых операций, дождаться всех, и если хоть
одна упала — вернуть её ошибку, а остальных отменить. Руками это WaitGroup плюс
context плюс аккуратный сбор первой ошибки — много мелочи, в которой легко
ошибиться. golang.org/x/sync/errgroup связывает всё это в один объект.
g, ctx := errgroup.WithContext(ctx)for _, u := range urls { g.Go(func() error { return fetch(ctx, u) })}if err := g.Wait(); err != nil { // первая ненулевая ошибка уже отменила ctx для остальных задач return err}
g.Wait() блокируется до завершения всех задач и возвращает первую ошибку. А
производный ctx отменяется в момент этой первой ошибки — поэтому остальные
fetch, если они слушают ctx.Done(), свернутся сами, не доделывая ненужную
работу. Это пул-на-минималках с правильной обработкой ошибок из коробки.
Все примеры выше в браузере исполняются однопоточно: yaegi гоняет горутины на
одном ядре, по очереди. Этого хватает, чтобы увидеть логику — кто кого закрывает,
кто когда выходит. Но это не показывает настоящий параллелизм и не ловит гонки.
Вот характерный баг, который в браузере мог бы «случайно работать», а на реальном
многоядерном железе под -race падает сразу:
// ОПАСНО: общий счётчик без синхронизации.// На нескольких ядрах counter++ — это read-modify-write,// и инкременты воркеров затирают друг друга.var counter intvar wg sync.WaitGroupfor i := 0; i < 1000; i++ { wg.Add(1) go func() { defer wg.Done() counter++ // ГОНКА: одновременный доступ из многих горутин }()}wg.Wait()fmt.Println(counter) // почти наверняка < 1000, и каждый запуск разный
Запускать это в playground бесполезно: однопоточный интерпретатор инкременты не
перемешает, и вы увидите ровно 1000 — ложное «всё работает». На проде go test -race мгновенно укажет на конфликт. Чинится либо atomic.AddInt64, либо
sync.Mutex, либо вообще перестройкой так, чтобы счётчик жил в одной горутине, а
остальные слали ему значения по каналу — подробности в главах про
sync-примитивы и гонки.
Мораль: корректность конкурентного кода нельзя проверить глазами на одном
прогоне. Запускайте тесты под -race на настоящем железе.
Двойное закрытие или закрытие не тем. У канала ровно один владелец закрытия,
и это отправитель. Если отправителей несколько — закрывает отдельная горутина
после wg.Wait(), а не каждый воркер. Закрытие из читателя или из нескольких
мест — паника.
Отправка без ветки отмены. Любая out <- v в долгоживущий канал должна
стоять в select рядом с <-ctx.Done(). Иначе ушедший потребитель оставляет
висящую навсегда горутину — утечку.
Запустил горутину и забыл, кто её остановит. На каждую go func() должен
быть ответ: что и когда её завершит. Нет ответа — потенциальная утечка по
определению.
Проверка конкурентности глазами в playground. Однопоточный интерпретатор не
перемешивает горутины и не видит гонок. «Работает в браузере» ничего не говорит
о проде — только go test -race на многоядерном железе.
Паттерны — это типовые графы потоков данных: узлы-горутины, рёбра-каналы.
Корректность каждого сводится к двум вопросам, с которых начиналась глава: как
сигнал конца доходит до каждого узла и кто закрывает каждое ребро. Есть
чёткий ответ на оба — структура не утечёт.
Это центр топика 3 (9–12) (пайплайны, fan-in/fan-out) и топика 5 (16–20)
(пулы, семафоры, батчинг). В топике 7 (23–32) паттерны собираются в
полноценные сервисы поверх context и
sync-примитивов. Тесты задач прицельно проверяют
завершаемость и отсутствие утечек — ровно то, ради чего
во всех примерах выше стоит ветка <-ctx.Done() и единый владелец закрытия.
За границы каналов — в select (тайм-ауты, приоритеты) и
context (отмена и дедлайны, которые пронизывают весь граф).