Ограничение конкурентности в Go: errgroup, пулы воркеров и backpressure
Ограничивайте объём незавершённых работ, а не просто отменяйте их.
Неограниченный fan-out — это режим отказа по умолчанию в конкурентной среде Go: каждая инструкция go — это ещё одна задача в работе, без ограничений и без владельца. Эта статья о том, как дать им и то, и другое.
Отмена и ограничение — это отдельные механизмы управления. Отмена решает, когда работа прекращается; ограничение решает, сколько её выполняется одновременно. Сервис может идеально отменять операции, но всё равно упасть, потому что десять тысяч горутин каждая независимо решили, что сейчас подходящий момент позвать в базу данных.

Страницы вокруг этой охватывают другую половину проблемы. Go context.Context Done Right рассматривает контекст как плоскость управления для остановки работы; здесь же тема — плоскость ёмкости: errgroup.SetLimit, семафорные каналы, пулы воркеров и ограниченные очереди — плюс одна форма утечки, которая переживает errgroup.Wait, даже когда все пути отмены настроены.
Пять примитивов, одна таблица решений
| Примитив | Частичные результаты | Фейл-фаст | Обратное давление | Динамическая работа | Лучший выбор для |
|---|---|---|---|---|---|
sync.WaitGroup + errors.Join |
да | нет | нет | нет | партии, где важны все результаты |
errgroup.WithContext |
нет | да | нет | нет | fan-out, где первая ошибка завершает запрос |
errgroup.SetLimit(n) |
нет | да | да (блокирует Go) |
нет | ограниченный фейл-фаст fan-out |
| Семафорный канал | да | вручную | да (блокирует захват) | да | пользовательский контроль доступа |
| Пул воркеров | да | вручную | да (очередь заполняется) | да | непрерывные потоки, долгоживущие воркеры |
Два вопроса выбирают строку. Первый: если одна задача падает, нужны ли результаты остальных? Если да, семантика групп фейл-фаст отбрасывает необходимую работу. Второй: работа приходит непрерывно или это фиксированная партия, которую можно посчитать перед запуском? Группы работают с партиями; пулы и очереди — с потоками.
errors.Join делает строку с обычным WaitGroup жизнеспособной — он собирает ошибки всех воркеров, а не хранит только первую, что хорошо сочетается с правилами перевода на границах в Go Error Handling Architecture.
Утечка, переживающая errgroup.Wait
Самая известная ошибка использования errgroup не имеет отношения к ограничению. Это разрыв в отмене, который прячется за кодом, который выглядит корректно:
g, gctx := errgroup.WithContext(ctx)
results := make(chan int) // без буфера
go func() { // производитель — никто не владеет этой горутиной
for i := 0; i < 100; i++ {
results <- i
}
close(results)
}()
g.Go(func() error { // потребитель — член группы
for r := range results {
if r == 5 {
return fmt.Errorf("save failed")
}
}
return nil
})
err := g.Wait() // возвращается при ошибке потребителя
Когда потребитель возвращает управление, gctx отменяется, и Wait возвращается. Производитель застрял на results <- i — и отмена контекста не разблокирует отправку по каналу. Никто не владеет этой горутиной, поэтому она остаётся застрявшей на всё время жизни процесса. Запустите эту конструкцию в цикле, и утечка будет строго линейной:
| Вариант | Живых горутин после 500 итераций |
|---|---|
Отправка без ветки ctx.Done |
501 (+500 утечек, по одной на итерацию) |
Отправка в select с <-gctx.Done() |
502 (+1 базовая) |
// исправление: каждая операция с каналом внутри отменяемой задачи
// получает ветку Done
select {
case results <- i:
case <-gctx.Done():
return
}
Та же самая разрыв существует и на стороне приёма. Правило механическое: каждая операция с каналом внутри кода, который может быть отменён, получает ветку <-ctx.Done(), и каждая горутина имеет владельца, который ждёт её или отменяет. goleak в CI ловит случаи, которые упускает ревью, — он встраивается в тот же автоматический контроль качества, что и инструменты из Go Linters: Essential Tools for Code Quality.
Две связанные формы стоит назвать. Если производитель является членом группы, тот же код приводит не к утечке, а к дедлоку — Wait блокируется навсегда, ожидая застрявшей отправки, что хотя бы громко падает. А производитель, который отправляет через select, но принимает из канала, чей потребитель уже ушёл, течёт тем же самым образом; ветки Done нужны с обеих сторон.
errgroup.SetLimit: ограниченный фейл-фаст fan-out
SetLimit превращает errgroup в группу с контролем допуска. Каждый вызов Go за пределами лимита блокируется, пока не освободится слот:
g, ctx := errgroup.WithContext(ctx)
g.SetLimit(8)
for _, item := range items {
g.Go(func() error {
return process(ctx, item)
})
}
err := g.Wait()
Три семантики важны в продакшене. Go блокируется, поэтому цикл выше применяет обратное давление к производителю — обычно это то, что нужно, но это значит, что горутина вызывающего может зависнуть, поэтому обработчик запросов, подающий большую партию, нуждается в собственном бюджете тайм-аута вокруг всего цикла. SetLimit(0) блокирует каждый вызов Go навсегда. И изменение лимита при активных горутинах приводит к панике — лимит фиксирован на всё время жизни группы.
SetLimit наследует семантику групп: первая ошибка отменяет контекст, и Wait возвращает первую ошибку, отбрасывая результаты работы, которая ещё в полёте. Это правильный примитив, когда падение всей операции при первом сбое — это правильно: получение одной страницы данных, вызов набора независимых health-check’ов, рассылка запроса на шарды.
Семафорные каналы: обратное давление как функция
Буферизованный канал ёмкости n — это семафор без зависимостей:
sem := make(chan struct{}, 8)
var wg sync.WaitGroup
for _, item := range items {
wg.Add(1)
sem <- struct{}{} // блокируется, когда 8 задач в работе
go func() {
defer wg.Done()
defer func() { <-sem }()
_ = process(ctx, item)
}()
}
wg.Wait()
Захват происходит до инструкции go, поэтому горутин создаются только когда есть слот. Размер буфера — это контракт: он одновременно является потолком конкурентности и очередью ожидающих допусков.
Примитив оправдывает своё место, когда нужен контроль доступа, который группы не выражают: захват с учётом контекста, пулы по отдельным downstream’ам или частичные результаты при сбое:
select {
case sem <- struct{}{}: // слот получен
case <-ctx.Done(): // вызывающий сдался, не получив доступа
return ctx.Err()
}
Измерения на одной и той же нагрузке показывают, что семафорный канал и SetLimit(8) финишируют с разницей менее 1,5% — выбор между ними основан на семантике, а не на производительности.
Пулы воркеров: для потоков, не для партий
Фиксированный пул долгоживущих воркеров разделяет допуск (очередь) и выполнение (воркеры):
func pool(ctx context.Context, workers int) chan<- func() {
jobs := make(chan func())
var wg sync.WaitGroup
for i := 0; i < workers; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for {
select {
case f, ok := <-jobs:
if !ok {
return
}
f()
case <-ctx.Done():
return
}
}
}()
}
return jobs
}
Пулы подходят для работы, которая приходит непрерывно — потребители очередей, циклы опроса, воркеры запросов, — где группа для каждой партии создавалась бы и разрушалась бесконечно. Начиная с Go 1.25, sync.WaitGroup.Go убирает шаблонный код Add/Done для распространённых случаев. Контракт завершения — это то, что нужно сделать правильно: закрытие jobs завершает воркеров после опустошения очереди; отмена ctx оставляет работу в очереди без внимания. Оба варианта легитимны, но только один может быть задокументированным поведением.
Ограниченные очереди: блокировать или отбрасывать
Буферизованный канал — это также очередь с потолком. Когда буфер заполняется, производитель блокируется — и этот момент является проекционным решением. Блокирование передаёт давление вверх по потоку тому, кто поставляет в очередь; отбрасывание сбрасывает нагрузку ценой потери работы. Мониторинговый конвейер может отбрасывать; конвейер заказов должен блокироваться или переливаться в долговременное хранилище.
select {
case queue <- event: // допущено
default: // очередь заполнена: политика отбрасывания здесь
dropped.Add(1)
}
Какова бы ни была политика, делайте её явной и измеряемой. Молчаливая ветка default — это то, куда события уходят, чтобы исчезнуть.
Стоимое ограничение — измеренное
Та же нагрузка, 10 000 задач по 5 мс имитации I/O, Go 1.27.1:
| Стратегия | Пик горутин | Стендовое время |
|---|---|---|
Неограниченный go на задачу |
~10 000 запущено | 18 мс |
errgroup.SetLimit(8) |
8 | 6,9 с |
| Семафорный канал (ёмкость 8) | 8 | 6,6 с |
Неограниченный прогон выигрывает по стендовому времени, потому что в эксперименте ничего не сопротивляется — десять тысяч таймерных горутин просто спят параллельно. Это и есть ловушка. Цена неограниченного fan-out никогда не в CPU вашего процесса; это десять тысяч одновременных соединений, downstream, который начинает тайм-аутоваться под ударом, и график памяти, который следует за глубиной очереди того, кого вы вызываете. На реальной зависимости неограниченный прогон не финиширует за 18 мс — он тайм-аутится. Ограниченные прогоны платят 6,6 секунд, потому что 10 000 задач ÷ 8 слотов × 5 мс — это 6,25 с неизбежной сериализации, и потолок в этом и состоит: он превращает неконтролируемый удар в предсказуемый сток за 6,25 секунды.
Одно уточнение важно, когда потолок выше, чем горстка слотов: держите потолок на уровне или ниже того, что downstream может устойчиво обрабатывать параллельно, и выводите его из этого лимита (размер пула соединений, rate limit, ёмкость воркеров), а не из ощущения о «разумной» параллельности.
Измерение и тестирование ограниченного кода
Три сигнала покрывают большую часть. Тенденции числа горутин (/sched/goroutines:goroutines, экспортируемые через Prometheus Go collector) ловят утечки как наклон. Задержка планировщика (/sched/latencies:seconds) ловит насыщение — исполняемые горутин, ожидающие CPU, это реальное давление, независимо от числа. Дельта pprof двух дампов горутин точно указывает, где живут застрявшие горутин.
Для тестов ограниченные воркеры — это тот случай, для которого создавался Testing Concurrent Go Code with synctest: псевдо-время позволяет полный сток 10 000 задач пройти за миллисекунды, synctest.Wait заменяет «поспал и надеюсь» для умиротворения, а пузырь быстро падает, если воркер остаётся устойчиво заблокированным — репро утечки выше умирает в тестовом пузыре, а не в продакшене.
На уровне сервиса тот же вопрос ограничения возникает между сервисами, где семафор становится rate limit, а circuit breaker защищает downstream — Circuit Breaker Pattern in Go охватывает этот слой, а Go Microservices for AI/ML Orchestration показывает очередевые формы, которые принимает слой оркестрации.
Выбор
при первом сбое?} D -- да --> E[WaitGroup + errors.Join
или семафорный канал] D -- нет --> F[errgroup.WithContext] F --> G{Нужен ли потолок конкурентности?} G -- да --> H[errgroup.SetLimit] G -- нет --> I[Группа без лимита
только когда партия доказуемо мала] C --> J{Очередь полная: блокировать или отбрасывать?} J -- блокировать --> K[Обратное давление вверх] J -- отбрасывать --> L[Сбросить нагрузку, посчитать отброшенные]
Таблица и диаграмма сжимаются в одну привычку: перед тем как писать go, назовите потолок и владельца. Потолок — это SetLimit, семафор или ограниченная очередь; владелец — это Wait, ветка ctx.Done на каждой операции с каналом, или и то, и другое. Каждая точка создания горутин, которая может ответить на эти два имени с одного взгляда, не разбудит никого в 3 ночи.