Ограничение конкурентности в 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 показывает очередевые формы, которые принимает слой оркестрации.

Выбор

flowchart TD A[Работа для конкурентного выполнения] --> B{Фиксированная партия или непрерывный поток?} B -- непрерывный поток --> C[Пул воркеров или ограниченная очередь] B -- фиксированная партия --> D{Нужны ли частичные результаты
при первом сбое?} 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 ночи.

Полезные ссылки

Подписаться

Получайте новые материалы про системы, инфраструктуру и AI engineering.