Отдельные примитивы — горутины, каналы,
select, context — это кирпичи. Сами по себе они
мало что значат: интересное начинается, когда из них складывают устойчивые
конструкции, повторяющиеся в любом проде. Пайплайн, который тянет данные через
несколько стадий. Пул воркеров, который держит нагрузку в узде. Семафор, который
не даёт открыть тысячу соединений разом. Батчер, который копит мелочь и пишет в
базу пачками.
Эта глава — каталог таких паттернов. Но не «вот код, копируй». Главная мысль
сквозная и простая: любой конкурентный граф корректен ровно настолько,
насколько чётко в нём отвечены два вопроса — как до каждого узла доходит сигнал
«данные кончились / пора останавливаться» и кто закрывает каждый канал. Ответы задают протокол завершения, который нужно выполнить на каждом блокирующем пути. Затем проверьте, что вызываемые операции возвращаются или реагируют на отмену; иначе возможна утечка даже при правильном закрытии каналов.
Примеры ниже по возможности запускаемые прямо в браузере. Там, где речь о
настоящем параллелизме или гонках, — статичный код с разбором, потому что в
браузере всё исполняется однопоточно и race-детектора нет.
Нужны каналы, select, WaitGroup, срезы и context. Начните с полного прохождения маленького pipeline и фиксированного пула, затем рассмотрите уход потребителя раньше конца. У каждого ожидания проверьте обычный выход и путь отмены; одного ответа «канал закрывает отправитель» для всего графа недостаточно.
На первом проходе изучите pipeline, fan-in/fan-out и worker pool. Семафор, batching и errgroup — следующие шаги. Алгоритмические мосты в конце главы готовят к задачам про высокую нагрузку, порядок и масштабирование. Backpressure означает обратное давление: если принимающая сторона не успевает, очередь заполняется или отправитель ждёт, ограничивая скорость источника.
Пайплайн — это конвейер. Цепочка стадий, где выход одной становится входом
следующей. Каждая стадия — горутина: читает из входного канала, что-то делает,
пишет в выходной. Стадии работают одновременно, и пока последняя печатает первый
результат, первая уже готовит десятый.
Ключевая дисциплина одна: стадия закрывает свой выходной канал, когда её вход
иссяк. Обычно через 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 { select { case <-ctx.Done(): return // отмена прерывает и ожидание входа case n, ok := <-in: if !ok { return } select { case out <- n * n: case <-ctx.Done(): return } } } }() return out}
Ветка <-ctx.Done() позволяет выйти и из ожидания входа, и из ожидания получателя. defer close(out) сообщает следующей стадии, что выход исчерпан. Отмена не передаётся через этот close: общий ctx нужно передать всем блокирующим стадиям, включая источник. gen из первого примера не принимает ctx, поэтому одна отменяемая sq ещё не обеспечивает остановку всего графа. Проверяйте оба ожидания каждой стадии и кооперативное завершение вызываемых функций.
Иногда одна стадия — узкое место: данные простые, а их много. Тогда её
распараллеливают.
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(), свернутся сами, не доделывая ненужную
работу. Это пул-на-минималках с правильной обработкой ошибок из коробки.
RPS — число начатых запросов в секунду. Ограничение workers = 3 означает не более трёх одновременно занятых работников, но не ограничивает RPS: быстрые вызовы успеют выполнить много запросов за секунду. Общий тикер перед началом вызова ограничивает темп всего пула; отдельный тикер у каждого работника умножил бы суммарный темп. Backlog — уже принятые задачи, ожидающие выполнения.
Token bucket хранит ограниченный запас разрешений: запрос тратит одно, время пополняет запас с заданным темпом до ёмкости. Поэтому разрешён короткий всплеск, после которого выдача восстанавливается постепенно. Leaky bucket выдаёт разрешения равномерно и заставляет ожидающих ждать. Для трёх токенов и скорости три токена в секунду начальные три операции можно разрешить сразу; при дробном учёте накопления следующий полный токен появится примерно через треть секунды. Это отличается от сброса счётчика раз в фиксированное окно.
Singleflight объединяет одновременные запросы одного ключа в одну выполняющуюся операцию и раздаёт её результат ожидающим. После завершения новой волне разрешено выполнить операцию снова. Кэш хранит результат дольше самой операции. Cache stampede — множество одинаковых запросов к источнику, возникших после исчезновения популярного значения из кэша; обычная защита map мьютексом ещё не объединяет эти запросы.
Шардирование разбивает карту на несколько независимых частей с собственными блокировками. Хеш переводит ключ в число; остаток hash % 32 выбирает одну из 32 частей. Один ключ должен всегда попадать в одну часть. Сначала проверьте этот инвариант для нескольких строк; затем добавляйте конкурентные обращения. Распределение ключей и более быстрый код — разные проверки: производительность измеряют отдельно.
В графе пользователи — вершины, дружба — связи. Обход в ширину (BFS) сначала проверяет ближайших соседей, затем соседей следующего уровня. Для графа 1 → [2, 3], 2 → [4], 3 → [4, 5], 4 → [1] выпишите новые вершины по числу переходов от 1: один переход — 2 и 3, два — 4 и 5; новых вершин на третьем переходе нет. visited — множество уже обнаруженных вершин: оно предотвращает повторный обход 4 и цикл обратно к 1. Начальную вершину помечают отдельно согласно контракту результата.
До горутин реализуйте последовательный обход: текущий уровень хранится в срезе, найденные новые соседи попадают в следующий. При конкурентном обходе обработайте вершины одного уровня параллельно и дождитесь всего уровня перед переходом к следующему. Проверку visited и добавление в него выполняйте в одной критической секции, чтобы два работника не взяли одну вершину. Семафор ограничивает одновременные вызовы источника, а не число найденных пользователей.
Для приоритетов 1..3 достаточно трёх очередей-срезов. В общей критической секции добавьте задачу в нужную очередь. Воркер выбирает первую задачу из очереди 3, иначе из 2, иначе из 1; если все пусты, ждёт Cond в цикле проверки. Для входа (1, A), (3, B), (2, C) выдача из уже наполненной очереди — B, C, A. Приоритет относится к выбору следующей задачи, а не к вытеснению уже выполняющейся.
container/heap хранит элементы так, чтобы быстро извлекать наивысший приоритет. После вставки и удаления он переставляет элементы, восстанавливая это правило; эти шаги называют sift-up и sift-down. Если устройство heap ещё незнакомо, для трёх фиксированных уровней начните с очередей. Строгий приоритет может оставить низкий уровень ждать при непрерывном потоке высокого; требование справедливости нуждается в отдельной политике, например квоте.
Результат второй задачи может быть готов раньше первого. Чтобы сохранить порядок, нужно отдельно учитывать готовность и очередь выдачи. Один путь — нумеровать входы и хранить готовые результаты по индексу, выдавая только следующий ожидаемый индекс.
Другой путь — канал каналов: chan (<-chan int) переносит каналы, из которых можно читать int. Для каждой задачи создают отдельный канал результата и кладут его в очередь в порядке входа. Работник может наполнить второй канал раньше первого, но коллектор сначала читает первый, затем второй. Если первый результат задержался, последующие готовы, однако пока не выдаются: это цена требования порядка. Перед большим пулом разберите две задачи и два таких канала вручную. Канал результата с буфером 1 позволяет работнику закончить отправку, пока коллектор ждёт предыдущий.
Перед задачами 23 и 32 реализуйте фиксированный пул. Затем отделите счётчик живых воркеров от числа занятых задач: спящий воркер тоже жив. Менеджер запускает дополнительного работника только если есть нагрузка и лимит не достигнут; решение и изменение счётчика должны быть согласованы. Уход лишнего воркера не должен уменьшать работающий пул ниже минимума. После Stop минимум уже не действует: все воркеры должны завершиться.
Дренировать очередь означает выполнить все ранее принятые задачи перед выходом. Остановить приём означает отказать новым; отменить выполняющиеся означает попросить их остановиться; дождаться означает подтвердить, что они вышли. Это отдельные действия. Произвольную func() нельзя безопасно принудительно оборвать внутри процесса; callback с контекстом должен проверять отмену сам. Возврат Stop по таймауту сообщает, что время ожидания исчерпано, и сам по себе не завершает медленную функцию.
В периодическом планировщике Ticker отвечает за «когда передать задачу», пул — за «кто её выполнит». ticker.Stop() не закрывает ticker.C: ожидание нужно совмещать с done/ctx. Паника из задачи требует recover в отдельном вызове одной задачи; если поставить его только на внешнюю worker-функцию, после восстановления сама worker-функция вернётся и не возьмёт следующую задачу.
Pub/sub отличается от fan-out работников: сообщение доставляется каждому подходящему подписчику, а не одному из конкурирующих получателей. При маршрутизации по тегам сначала получите список подходящих обработчиков под защитой карты, затем вызывайте их с учётом выбранного контракта. Вызов произвольного обработчика под общим мьютексом может задержать подписку, публикацию и остановку всего маршрутизатора.
Все примеры выше в браузере исполняются однопоточно: 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 на многоядерном железе.
В простом pipeline gen → sq потребитель принимает только первое число и уходит. Назовите операции, на которых оставшиеся горутины могут ждать. Затем сделайте обе стадии отменяемыми: отдельно на ожидании входа и отправке выхода. Передайте им общий контекст, после первого результата отмените его и через WaitGroup дождитесь завершения всех стадий. Не пытайтесь исправить всё только большим буфером.
После ухода потребителя sq может ждать отправки следующего квадрата, а gen — передачи следующего числа в sq. Оба ожидания должны слышать общий сигнал отмены. Каркас одной стадии:
func squareStage(ctx context.Context, in <-chan int, wg *sync.WaitGroup) <-chan int { out := make(chan int) wg.Add(1) go func() { defer wg.Done() defer close(out) for { select { case <-ctx.Done(): return case n, ok := <-in: if !ok { return } select { case out <- n * n: case <-ctx.Done(): return } } } }() return out}
Источник должен использовать такой же select при каждой отправке и выполнить Done при выходе. Зарегистрируйте обе стадии до вызова Wait. Если у входа нет данных, отменяется первый select; если потребитель не читает результат, второй. Close передаёт конец вниз, общий ctx сообщает отмену всем стадиям; для прекращения источника недостаточно только закрыть выход последней стадии.
Паттерны — это типовые графы потоков данных: узлы-горутины, рёбра-каналы.
Корректность каждого сводится к двум вопросам, с которых начиналась глава: как
сигнал конца доходит до каждого узла и кто закрывает каждое ребро. Ответы нужны для протокола завершения; отсутствие утечек дополнительно требует отменяемых ожиданий, возвращающихся обработчиков и проверки всех путей.
Это центр топика 3 (9–12) (пайплайны, fan-in/fan-out) и топика 5 (16–20)
(пулы, семафоры, батчинг). В топике 7 (23–32) паттерны собираются в
полноценные сервисы поверх context и
sync-примитивов. Тесты задач прицельно проверяют
завершаемость и отсутствие утечек — ровно то, ради чего
во всех примерах выше стоит ветка <-ctx.Done() и единый владелец закрытия.
За границы каналов — в select (тайм-ауты, приоритеты) и
context (отмена и дедлайны, которые пронизывают весь граф).