GraphLMS

Go
0/32 решеноНачать
Глава 08 · Основы~15 мин чтения

Паттерны конкурентности

О чём глава

Отдельные примитивы — горутины, каналы, select, context — это кирпичи. Сами по себе они мало что значат: интересное начинается, когда из них складывают устойчивые конструкции, повторяющиеся в любом проде. Пайплайн, который тянет данные через несколько стадий. Пул воркеров, который держит нагрузку в узде. Семафор, который не даёт открыть тысячу соединений разом. Батчер, который копит мелочь и пишет в базу пачками.

Эта глава — каталог таких паттернов. Но не «вот код, копируй». Главная мысль сквозная и простая: любой конкурентный граф корректен ровно настолько, насколько чётко в нём отвечены два вопроса — как до каждого узла доходит сигнал «данные кончились / пора останавливаться» и кто закрывает каждый канал. Если на оба есть ответ — структура не утечёт. Если хоть на один ответа нет — рано или поздно повиснет горутина.

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

Pipeline

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

Ключевая дисциплина одна: стадия закрывает свой выходной канал, когда её вход иссяк. Обычно через defer close(out). Это и есть способ передать «данные кончились» вниз по конвейеру — следующая стадия увидит закрытие через range и сама закроется. Сигнал конца течёт по тем же рельсам, что и данные.

package main
 
import "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-out — несколько горутин читают из одного входного канала. Каждая берёт следующий элемент, как только освободилась; быстрые забирают больше, медленные меньше — нагрузка балансируется сама. Fan-in — обратная операция: слить несколько выходных каналов обратно в один, чтобы потребитель видел единый поток.

Fan-in — это та самая задача «несколько отправителей, один канал» из главы про каналы. И тонкость там ровно одна: выходной канал закрывает кто-то один, и только когда замолчали все источники. Иначе либо запись в закрытый канал (паника), либо double-close (тоже паника). Идиома — sync.WaitGroup: каждый источник делает Done, отдельная горутина ждёт Wait и закрывает выход.

package main
 
import (
	"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: слить каналы в один, закрыть после Wait
func 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 — по той же причине, что и в пайплайне: чтобы источники не зависли, если потребитель ушёл раньше.

Worker pool

Fan-out с каналом квадратов — это уже почти пул. Worker pool делает идею явной: фиксированное число воркеров разбирает задачи из общей очереди и складывает результаты в общий канал. Зачем фиксированное? Потому что «горутина на задачу» при миллионе задач — это миллион одновременных обращений к базе или диску. Пул из восьми воркеров обрабатывает всё то же, но держит конкуренцию ровно на восьми. Размер пула — ваш главный рычаг backpressure.

package main
 
import (
	"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. Воркеров много, владелец закрытия — один.

Semaphore (ограничение конкуренции)

Пул хорош, когда есть поток однотипных задач и канал результатов. Но иногда нужно проще: «запусти вот эти задачи, но не больше N одновременно». Заводить ради этого пул и канал результатов — перебор. Хватает буферизованного канала как счётного семафора.

Идея игрушечная: канал на N элементов — это N пропусков. Чтобы стартовать, горутина кладёт пропуск в канал (заняла слот). Если все N слотов заняты — отправка блокируется, и новая горутина ждёт. Закончила — забирает пропуск обратно (освободила слот).

package main
 
import (
	"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, чтобы ожидание слота было отменяемым.

Batching

Последний паттерн — про экономию. Писать в базу по одной записи дорого: каждая — round-trip. Батчинг копит элементы и сбрасывает их пачкой. Но если просто ждать, пока наберётся size, можно зависнуть: пришло три элемента, поток затих — и они никогда не запишутся. Поэтому у батчера два триггера: набралось size или вышло время, что наступит раньше. Оба живут в одном select.

package main
 
import (
	"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() — иначе утечёт сам таймер.

errgroup: оркестрация с ошибкой

Очень частая задача: запустить 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 int
var wg sync.WaitGroup
for 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 (отмена и дедлайны, которые пронизывают весь граф).