Паттерны конкурентности в Go: pipeline, semaphore, errgroup, singleflight

Рабочие паттерны конкурентности в Go: pipeline с отменой и очисткой, fan-in/fan-out, worker pool с ограничением параллелизма, semaphore, errgroup, singleflight и батчинг — и когда паттерн лишний

Горутины и каналы — это кирпичи; из них складываются повторяющиеся конструкции, которые встречаются в почти любом сервисе: конвейер обработки, пул воркеров с ограничением, дедупликация одинаковых запросов. Знание этих паттернов экономит время и, что важнее, помогает не наделать классических ошибок — утечек горутин, незакрытых каналов и неконтролируемого параллелизма.

Это четвёртая статья серии «Конкурентность в Go». Она опирается на основы из первой и синхронизацию из второй; модель памяти и правила видимости разобраны в третьей, а как ловить утечки и гонки, которые проскочат в этих паттернах, — в пятой.

Статья код-центричная: под каждый паттерн — минимальный, но корректный и компилирующийся пример, с явным разбором, когда паттерн применим и как из него не утечь горутинами.

Мастерская паттернов конкурентности Go: pipeline (стадии decode/validate/enrich/store), fan-out/fan-in, worker pool с семафором (MAX=3), errgroup с кнопкой CTX CANCEL — одна ошибка останавливает всех воркеров

В статье

Общий принцип: отмена и очистка

Все паттерны ниже держатся на двух правилах, нарушение которых и есть источник утечек.

  1. У каждой запущенной горутины должен быть однозначный путь завершения. Если горутина блокируется на отправке в канал, который никто больше не читает, или на чтении из канала, который никто не закроет, — это утечка. Она не всегда фатальна сразу, но накапливается под нагрузкой.
  2. Кто создал канал — тот и закрывает его, и делает это ровно один раз. Стадия-владелец закрывает свой выходной канал через defer close(out), когда её цикл завершился. Закрытие сигнализирует потребителям «данных больше не будет», и их range корректно останавливается.

Инструмент отмены — context.Context. Он пробрасывается через все стадии, и в каждой точке, где горутина может заблокироваться на отправке, стоит select с веткой <-ctx.Done(). Без этой ветки ранний выход потребителя (break, ошибка, таймаут) оставит верхние стадии висеть на out <- v навсегда.

Подробно про context — в блоге Go; ниже он используется как данность.

Pipeline: стадии на каналах с отменой

Конвейер — это цепочка стадий, где выход одной стадии (<-chan T) является входом следующей. Каждая стадия — функция, запускающая горутину, возвращающая канал и закрывающая его по завершении.

// gen — источник (generator): отдаёт числа в канал и закрывает его.
func gen(ctx context.Context, nums ...int) <-chan int {
	out := make(chan int)
	go func() {
		defer close(out) // владелец канала закрывает его сам
		for _, n := range nums {
			select {
			case out <- n:
			case <-ctx.Done(): // потребитель ушёл — не блокируемся навсегда
				return
			}
		}
	}()
	return out
}

// sq — стадия обработки: читает из in, пишет квадраты в out.
func sq(ctx context.Context, in <-chan int) <-chan int {
	out := make(chan int)
	go func() {
		defer close(out)
		for n := range in { // range завершится, когда in закроют
			select {
			case out <- n * n:
			case <-ctx.Done():
				return
			}
		}
	}()
	return out
}

Сборка и потребление. Ключевой момент — defer cancel(): он гарантирует, что даже при раннем выходе из цикла все верхние стадии получат сигнал отмены и завершатся.

ctx, cancel := context.WithCancel(context.Background())
defer cancel() // освобождает стадии, если выйдем из range раньше времени

for n := range sq(ctx, gen(ctx, 2, 3, 4, 5)) {
	fmt.Println(n)
	if n == 9 {
		break // ранний выход: gen и sq зависли бы на send без cancel()
	}
}

Почему это не течёт. При break потребитель перестаёт читать из sq. Горутина sq в этот момент может быть заблокирована на out <- n*n, а gen — на out <- n. Отложенный cancel() закрывает ctx.Done(), обе горутины выходят по второй ветке select, отрабатывают defer close(out) и завершаются. Без cancel() (или без ветки <-ctx.Done() в select) обе остались бы висеть.

Ошибки в конвейере. Канал значений несёт только данные; для ошибок есть два подхода. Простой — передавать по каналу структуру struct{ Val T; Err error } и разбирать её на приёме. Второй, когда ошибка любой стадии должна остановить весь конвейер, — построить пайплайн поверх errgroup: первая ошибка отменяет общий контекст, и все стадии сворачиваются через <-ctx.Done().

Канонический разбор конвейеров с отменой — Go blog: Pipelines and cancellation.

Fan-out / fan-in

Fan-out — распараллелить медленную стадию: несколько одинаковых горутин читают из одного входного канала. Пока одна занята, другие разбирают остальные элементы. Балансировка получается сама собой: каждый воркер берёт следующий элемент, когда освободился.

Fan-in — слить результаты нескольких каналов в один. Здесь нужен WaitGroup: закрыть общий выход можно только после того, как все источники исчерпаны.

// merge сливает несколько каналов в один и закрывает out,
// когда все входные каналы закрыты (или ctx отменён).
func merge(ctx context.Context, cs ...<-chan int) <-chan int {
	var wg sync.WaitGroup
	out := make(chan int)

	forward := func(c <-chan int) {
		defer wg.Done()
		for n := range c {
			select {
			case out <- n:
			case <-ctx.Done():
				return
			}
		}
	}

	wg.Add(len(cs))
	for _, c := range cs {
		go forward(c)
	}

	// Отдельная горутина закрывает out ровно один раз — после всех forward.
	go func() {
		wg.Wait()
		close(out)
	}()
	return out
}

Fan-out поверх конвейера из предыдущего раздела:

in := gen(ctx, 1, 2, 3, 4, 5, 6, 7, 8, 9)

// Три параллельных экземпляра стадии sq читают из одного in.
c1 := sq(ctx, in)
c2 := sq(ctx, in)
c3 := sq(ctx, in)

// Fan-in: собираем результаты в один поток.
for n := range merge(ctx, c1, c2, c3) {
	fmt.Println(n)
}

Почему это корректно. Три горутины sq конкурентно читают из in; когда gen закроет in, каждый range in завершится, каждый sq закроет свой выход. В merge три forward увидят закрытие своих c, отработают wg.Done(), и только тогда закроется out. Ни одна горутина не остаётся висеть. При раннем выходе потребителя merge — тот же приём с cancel(), что и в конвейере.

Worker pool и ограничение параллелизма

Fan-out запускает по горутине на источник; часто нужно ограничить число одновременно работающих горутин фиксированной величиной — чтобы не перегрузить БД, внешний API или диск. Два основных способа.

Пул фиксированного размера

N заранее запущенных воркеров читают задачи из общего канала. Число горутин постоянно и не зависит от числа задач.

func workerPool(ctx context.Context, tasks <-chan Task, workers int) <-chan Result {
	results := make(chan Result)
	var wg sync.WaitGroup

	wg.Add(workers)
	for i := 0; i < workers; i++ {
		go func() {
			defer wg.Done()
			for t := range tasks { // все воркеры делят один канал задач
				select {
				case results <- process(ctx, t):
				case <-ctx.Done():
					return
				}
			}
		}()
	}

	go func() {
		wg.Wait()
		close(results)
	}()
	return results
}

Отправитель задач закрывает tasks, когда задачи кончились; воркеры выходят из range, wg доходит до нуля, results закрывается. Потребитель results обязан либо дочитать канал до конца, либо отменить ctx — иначе воркеры зависнут на results <- ....

Семафор на буферизованном канале

Когда воркеры не постоянные, а горутина создаётся под каждую задачу, число одновременных ограничивают буферизованным каналом-семафором: ёмкость буфера = максимум параллелизма.

func processAll(ctx context.Context, items []Item, limit int) {
	sem := make(chan struct{}, limit) // ёмкость = лимит параллелизма
	var wg sync.WaitGroup

	for _, item := range items {
		select {
		case sem <- struct{}{}: // занять слот; блокирует, когда все limit заняты
		case <-ctx.Done():
			wg.Wait()
			return
		}
		wg.Add(1)
		go func(item Item) {
			defer wg.Done()
			defer func() { <-sem }() // освободить слот — обязательно через defer
			handle(ctx, item)
		}(item)
	}
	wg.Wait()
}

Освобождение слота через defer критично: если handle запаникует, defer всё равно вернёт слот, и семафор не заклинит. item передаётся аргументом — привычка, не зависящая от версии; начиная с Go 1.22 переменная цикла и так своя на каждой итерации.

Подход Число горутин Когда выбирать
Пул фиксированного размера Постоянное (N) Долго живущий пул, поток задач через канал, важна предсказуемость
Семафор на канале По горутине на задачу, но одновременно ≤ limit Разовая пачка задач, разнородная работа, лаконичность
errgroup.SetLimit limit, управляет группа Нужны сбор ошибок и отмена по первой из них
semaphore.Weighted По горутине, но с весом Задачи разной «стоимости» (память, слоты)

Backpressure. Если воркеры хронически не успевают, задачи копятся. Неограниченный буфер перед пулом — это отложенная катастрофа по памяти. Здоровые варианты: небуферизованный вход (отправитель естественно тормозит, пока воркер не возьмёт задачу), ограниченная очередь с явным отказом при переполнении (load shedding) или сигнал вверх по стеку. Подробно — в статье Backpressure и load shedding.

Взвешенный semaphore для лимита ресурсов

Канал-семафор считает задачи штуками. Когда задачи потребляют ресурс неравномерно (одна тянет 1 ГБ, другая — 50 МБ), нужен взвешенный семафор — golang.org/x/sync/semaphore. Он ограничивает не число задач, а суммарный вес одновременно работающих.

import "golang.org/x/sync/semaphore"

// maxWeight — общий «бюджет» (например, число ядер или условных единиц памяти).
func runWeighted(ctx context.Context, tasks []Task, maxWeight int64) error {
	sem := semaphore.NewWeighted(maxWeight)
	var wg sync.WaitGroup

	for _, t := range tasks {
		// Acquire блокирует, пока не освободится t.Cost единиц,
		// и вернёт ошибку, если ctx отменён во время ожидания.
		if err := sem.Acquire(ctx, t.Cost); err != nil {
			break // прекращаем запуск новых задач
		}
		wg.Add(1)
		go func(t Task) {
			defer wg.Done()
			defer sem.Release(t.Cost) // вернуть ровно столько, сколько заняли
			process(ctx, t)
		}(t)
	}

	wg.Wait()
	return ctx.Err()
}

Acquire с весом больше maxWeight заблокируется навсегда — вес одной задачи не должен превышать бюджет. Release должен возвращать тот же вес, что был занят, иначе счётчик поедет. Для простого случая «не больше N одинаковых задач» взвешенный семафор избыточен — берите канал или errgroup.SetLimit.

errgroup: первая ошибка отменяет остальных

golang.org/x/sync/errgroup — самый практичный из паттернов для типовой задачи «запустить несколько параллельных вызовов зависимостей и дождаться их, а при первой ошибке отменить остальные». Он инкапсулирует WaitGroup + сбор первой ошибки + отмену контекста.

import "golang.org/x/sync/errgroup"

func fetchAll(ctx context.Context, urls []string) ([]Page, error) {
	// WithContext даёт производный ctx, который отменится при первой ошибке
	// любой из g.Go(...) или когда g.Wait() вернёт управление.
	g, ctx := errgroup.WithContext(ctx)
	pages := make([]Page, len(urls))

	for i, url := range urls {
		g.Go(func() error {
			p, err := fetch(ctx, url) // ctx уже «завязан» на группу
			if err != nil {
				return err // эта ошибка отменит ctx для остальных
			}
			pages[i] = p // каждая горутина пишет в свой индекс — гонки нет
			return nil
		})
	}

	if err := g.Wait(); err != nil {
		return nil, err // вернётся ПЕРВАЯ ненулевая ошибка
	}
	return pages, nil
}

Семантика, которую важно помнить:

  • g.Wait() возвращает первую ненулевую ошибку (последующие теряются) и ждёт завершения всех горутин.
  • контекст из WithContext отменяется, как только первая горутина вернула ошибку, — остальные обязаны реагировать на <-ctx.Done(), иначе отмена ничего не ускорит.
  • запись в общий срез по непересекающимся индексам безопасна без мьютекса; общая мапа или счётчик — уже нет, нужна синхронизация (см. sync и атомики).

Ограничение параллелизма — SetLimit. Метод SetLimit(n) (и неблокирующий TryGo) добавлены в golang.org/x/sync v0.1.0 (2022) и с тех пор — часть стабильного API. После SetLimit(n) вызов g.Go блокируется, пока активных горутин не станет меньше n. Это превращает errgroup в worker pool со сбором ошибок:

func fetchLimited(ctx context.Context, urls []string, limit int) error {
	g, ctx := errgroup.WithContext(ctx)
	g.SetLimit(limit) // не более limit одновременных fetch

	for _, url := range urls {
		g.Go(func() error {
			return fetch(ctx, url) // g.Go подождёт свободный слот
		})
	}
	return g.Wait()
}

SetLimit нужно вызвать до первого g.Go и нельзя менять, пока есть активные горутины. TryGo вернёт false вместо блокировки, если слот занят, — удобно, когда лишнюю задачу можно отбросить, а не ждать.

singleflight: схлопывание одинаковых запросов

golang.org/x/sync/singleflight решает проблему cache stampede: когда кэш протух и сотня одновременных запросов на один ключ бьёт в БД одновременно. singleflight выполняет реальный вызов один раз на группу одинаковых ключей, а остальным раздаёт тот же результат.

import "golang.org/x/sync/singleflight"

var group singleflight.Group

func getUser(ctx context.Context, id string) (*User, error) {
	// Все конкурентные Do с одинаковым id дождутся ОДНОГО вызова функции.
	v, err, shared := group.Do(id, func() (any, error) {
		return fetchUserFromDB(ctx, id)
	})
	if err != nil {
		return nil, err
	}
	_ = shared // true, если результат разделён с другими вызывающими
	return v.(*User), nil
}

Подводные камни, из-за которых singleflight применяют неправильно:

  • Ошибка тоже разделяется. Если единственный реальный вызов упал, ту же ошибку получат все ожидающие. Один сбойный запрос «заражает» всю группу — но только на время этого вызова, следующая волна запустит новый.
  • Отравленный кэш при связке с кэшем. Если внутри Do результат кладётся в кэш, а вызов вернул невалидное значение, оно раздастся всем. Вызывайте group.Forget(key), чтобы «забыть» ключ и позволить следующему запросу повторить вызов, а не присоединиться к уже провалившемуся.
  • Контекст первого вызывающего. Функция в Do выполняется в горутине первого вызвавшего. Если использовать внутри ctx конкретного вызова, его отмена (таймаут, разрыв соединения) убьёт общий вызов для всех ожидающих. Для честной отмены на вызывающего используйте DoChan + select:
func getUserCancelable(ctx context.Context, id string) (*User, error) {
	ch := group.DoChan(id, func() (any, error) {
		// ВАЖНО: не ctx конкретного вызывающего, а фоновый/производный —
		// чтобы отмена одного не оборвала общий вызов для остальных.
		return fetchUserFromDB(context.WithoutCancel(ctx), id)
	})

	select {
	case res := <-ch:
		if res.Err != nil {
			return nil, res.Err
		}
		return res.Val.(*User), nil
	case <-ctx.Done(): // отменяем только своё ожидание, вызов продолжается
		return nil, ctx.Err()
	}
}

context.WithoutCancel (Go 1.21+) отвязывает значения контекста от его отмены — общий вызов не оборвётся из-за одного ушедшего клиента. Так каждый вызывающий отменяет только собственное ожидание, а полезная работа доводится до конца для оставшихся.

Rate limiting: тикер и token bucket

Ограничение частоты (а не параллелизма) — отдельная задача. Два уровня.

Простой тикер — когда достаточно «не чаще, чем раз в T»:

func pollAtRate(ctx context.Context, work func(context.Context)) {
	tick := time.NewTicker(200 * time.Millisecond)
	defer tick.Stop() // Stop освобождает связанный с тикером таймер в рантайме
	for {
		select {
		case <-tick.C:
			work(ctx)
		case <-ctx.Done():
			return
		}
	}
}

Тикер не даёт всплесков: строго равномерный интервал. Забытый Stop() оставляет таймер тикера жить в рантайме — утечка ресурса (не отдельной горутины на тикер).

Token bucketgolang.org/x/time/rate — когда нужны средняя частота плюс допустимый всплеск (burst):

import "golang.org/x/time/rate"

func callAPI(ctx context.Context, reqs []Request) error {
	// 10 запросов/сек в среднем, но допускаем всплеск до 20 подряд.
	limiter := rate.NewLimiter(rate.Limit(10), 20)

	for _, r := range reqs {
		if err := limiter.Wait(ctx); err != nil {
			return err // ctx отменён во время ожидания токена
		}
		do(ctx, r)
	}
	return nil
}

Три режима под разные задачи:

  • Wait(ctx)блокирует до появления токена; естественный throttle для клиента.
  • Allow() — возвращает bool немедленно; удобно, чтобы отбросить запрос сверх лимита (server-side rate limit).
  • Reserve() — резервирует токен и сообщает задержку Delay(); для более тонкого контроля.

burst = 0 у Wait заблокирует навсегда (кроме rate.Inf) — токенов не будет никогда.

Батчинг и дебаунс

Батчинг собирает поток мелких элементов в пачки — чтобы амортизировать стоимость операции (один INSERT на 100 строк вместо 100 запросов). Условие сброса — по размеру ИЛИ по времени, что наступит раньше, иначе последние элементы залипнут в буфере.

func batch(ctx context.Context, in <-chan Item, maxSize int, maxWait time.Duration) <-chan []Item {
	out := make(chan []Item)
	go func() {
		defer close(out)

		var buf []Item
		timer := time.NewTimer(maxWait)
		timer.Stop() // стартуем без активного таймера
		defer timer.Stop()

		flush := func() {
			if len(buf) == 0 {
				return // защита от «пустого» сброса по устаревшему тику
			}
			select {
			case out <- buf:
			case <-ctx.Done():
			}
			buf = nil
			timer.Stop()
		}

		for {
			select {
			case item, ok := <-in:
				if !ok {
					flush() // входной канал закрыт — отдать остаток
					return
				}
				if len(buf) == 0 {
					timer.Reset(maxWait) // отсчёт с первого элемента пачки
				}
				buf = append(buf, item)
				if len(buf) >= maxSize {
					flush() // сброс по размеру
				}
			case <-timer.C:
				flush() // сброс по времени
			case <-ctx.Done():
				return
			}
		}
	}()
	return out
}

Тонкость с таймером: flush() защищён проверкой len(buf) == 0, поэтому даже устаревший тик из timer.C (возможный при гонке Stop/Reset до Go 1.23) приводит к безобидному no-op. Начиная с Go 1.23 семантика Timer/Ticker строже — устаревшие значения из канала больше не доставляются, — но код остаётся корректным на обеих версиях именно из-за этой проверки. Остаток при закрытии in отдаётся через финальный flush() — элементы не теряются.

Дебаунс — родственник: сбрасывать не пачку, а только последнее значение после паузы тишины (типично для «поиск по мере ввода»). Схема та же — таймер, сбрасываемый на каждом новом событии, только буфер хранит одно значение.

Вспомогательные каналы: or-done, tee, bridge

Три небольших паттерна из книги Кэтрин Кокс-Бьюдей «Concurrency in Go» убирают повторяющийся select-boilerplate. Использовать их стоит по необходимости, а не ради самих паттернов.

or-done-channel — инкапсулирует «читать канал, пока он не закрыт или ctx не отменён». Позволяет писать обычный range вместо вложенного select в местах потребления:

func orDone(ctx context.Context, in <-chan int) <-chan int {
	out := make(chan int)
	go func() {
		defer close(out)
		for {
			select {
			case <-ctx.Done():
				return
			case v, ok := <-in:
				if !ok {
					return
				}
				select {
				case out <- v:
				case <-ctx.Done():
					return
				}
			}
		}
	}()
	return out
}

// Использование: чистый range без ручного select на ctx.Done().
// for v := range orDone(ctx, someChannel) { use(v) }

tee-channel — раздваивает один вход на два независимых выхода (например, один поток на обработку, второй — на логирование). Оба выхода нужно читать, иначе отправитель заблокируется на медленном.

bridge-channel — превращает «канал каналов» (<-chan <-chan T) в один плоский поток, последовательно вычитывая каждый вложенный канал. Полезно, когда стадии сами порождают каналы.

Оба (tee, bridge) строятся на том же приёме — горутина-владелец, defer close, select на ctx.Done() при каждой отправке. Полные реализации есть в книге и в стенде к статье.

Когда паттерн — лишнее усложнение

Конкурентность — не самоцель, а инструмент под конкретное узкое место. Признаки, что паттерн здесь лишний:

  • Данных мало или работа быстрая. Конвейер из трёх горутин на срез из десяти элементов проиграет обычному циклу: накладные расходы на каналы и переключения превысят выигрыш.
  • Нужно просто защитить состояние. Счётчик, кэш, мапа под конкурентным доступом — это sync.Mutex или атомик, а не канал и горутина. Каналы координируют передачу владения, а не заменяют мьютекс (см. sync и атомики).
  • Стадии по своей природе последовательны. Если каждый следующий шаг зависит от результата предыдущего, распараллеливать нечего — конвейер не даст ускорения, только усложнит отмену и обработку ошибок.
  • Параллелизм не упирается в ресурс. Ограничивать параллелизм имеет смысл, когда есть что беречь (соединения к БД, память, внешний rate limit). Если узкого места нет — лишний семафор только добавляет кода.

Правило: сначала простой последовательный код, конкурентность — когда профиль (см. отладку гонок и утечек) показал реальную выгоду. Преждевременная конкурентность — источник тех самых утечек, ради которых написана эта статья.

Подводные камни (утечки в паттернах)

Почти все утечки горутин в этих паттернах сводятся к нескольким повторяющимся причинам.

  • Отправка в канал без ветки <-ctx.Done(). Потребитель вышел раньше — верхняя стадия навсегда виснет на out <- v. Каждый send в горутине-владельце должен быть в select с отменой.
  • Забыли закрыть выходной канал. range out у потребителя не завершится никогда — его горутина утечёт. Владелец обязан defer close(out).
  • Закрыли канал дважды или закрыл не владелец. close уже закрытого канала — паника. Закрывает ровно один и тот же владелец, ровно один раз.
  • errgroup: горутины игнорируют ctx. WithContext отменит контекст при первой ошибке, но если задачи не смотрят на <-ctx.Done(), отмена ничего не ускорит — группа дождётся всех.
  • Семафор без defer на освобождении. Паника или ранний return в задаче — слот не вернулся, семафор постепенно заклинивает. Освобождение всегда через defer.
  • Небуферизованный result-канал, который никто не дочитывает. Воркеры пула виснут на results <- .... Потребитель обязан либо дочитать канал, либо отменить ctx.
  • singleflight: отмена одного клиента рушит общий вызов. Использование ctx конкретного вызывающего внутри Do/DoChan — отмена одного обрывает работу для всех. Отвязывайте контекст (context.WithoutCancel) или используйте DoChan + select.
  • time.Ticker/time.Timer без Stop(). Утечка таймера. Всегда defer tick.Stop().
  • Утечка из-за раннего выхода без cancel(). context.WithCancel без отложенного cancel() — стадии не получат сигнал и повиснут. defer cancel() — обязателен.

Как эти утечки увидеть — через runtime.NumGoroutine, goleak, pprof-профиль горутин и -race — разобрано в пятой статье серии.

Checklist

  • у каждой запущенной горутины есть однозначный путь завершения;
  • каждый выходной канал закрывается его владельцем ровно один раз (defer close);
  • каждый send внутри горутины-владельца — в select с веткой <-ctx.Done();
  • context пробрасывается через все стадии, cancel вызывается через defer;
  • параллелизм ограничен там, где есть что беречь (БД, память, внешний лимит);
  • освобождение семафора/слота — через defer, устойчиво к панике;
  • errgroup-задачи реагируют на отмену контекста, запись в общий срез — по непересекающимся индексам;
  • singleflight не завязан на контекст одного вызывающего; при кэшировании продуман Forget;
  • Ticker/Timer останавливаются через Stop();
  • перед добавлением конкурентности проверено, что последовательный код действительно не подходит;
  • утечки проверены (goleak/pprof/-race) до вывода в прод.

Демо и версии

Запускаемые примеры всех паттернов — в живом стенде digital-cookbook/go-concurrency/patterns/: каждый паттерн отдельным файлом с тестом на отсутствие утечек (goleak), весь пакет проходит go test -race. Код в тексте — выжимки.

  • Go: фрагменты в тексте требуют Go 1.21+ (используют context.WithoutCancel из Go 1.21; на Go 1.23+ семантика таймеров строже, но код корректен и раньше — за счёт проверки len(buf) == 0), а стенд целиком объявляет Go 1.25+ (единый модуль серии с testing/synctest и WaitGroup.Go);
  • golang.org/x/sync — актуальная версия (errgroup.SetLimit/TryGo доступны с v0.1.0, 2022);
  • golang.org/x/time — для rate.Limiter.

Документация и первоисточники

Соседние статьи серии: горутины и каналы, sync и атомики, модель памяти и happens-before, отладка гонок и утечек. Смежное: graceful shutdown и backpressure и load shedding.

Обсуждение в Telegram

Присоединиться →

Комментарии