Go Конкурентность в Go: горутины, каналы, context, sync и паттерны
0%

Конкурентность в Go: горутины, каналы, context, sync и паттерны

Конкурентность в Go

Конкурентность — то, ради чего многие приходят в Go. В большинстве языков параллельное программирование — это дополнительная библиотека, ручное управление пулами потоков и минное поле блокировок. В Go конкурентность встроена в язык на уровне синтаксиса, а её философия умещается в одну фразу Роба Пайка: «Don’t communicate by sharing memory; share memory by communicating» — не разделяйте память, чтобы общаться; общайтесь, чтобы разделять память.

Сначала важное различие. Конкурентность (concurrency) — это про структуру: способ организовать программу как набор независимо продвигающихся задач. Параллелизм (parallelism) — про исполнение: одновременный запуск на нескольких ядрах. Конкурентная программа может исполняться и на одном ядре. Go даёт вам инструменты структурировать задачи конкурентно, а рантайм уже раскладывает их по ядрам параллельно. Обязательный доклад на эту тему — Concurrency is not Parallelism Роба Пайка.

Горутины

Горутина — это функция, запущенная конкурентно. Синтаксис — просто слово go перед вызовом:

go doWork()          // запустить doWork() конкурентно и сразу продолжить
go func() {          // или анонимную функцию
	fmt.Println("привет из горутины")
}()

Почему горутин можно иметь миллионы, а потоков ОС — тысячи? Потому что горутина стоит копейки: стартовый стек — около 2 КБ (растёт по мере надобности), тогда как поток ОС резервирует мегабайты. Планировщик Go (модель M:N из обзора) мультиплексирует горутины на небольшое число потоков. Когда горутина блокируется на сетевом I/O, планировщик снимает её и ставит другую — без дорогого переключения контекста в ядре.

Критически важная деталь: main не ждёт горутины. Когда main завершается, программа умирает вместе со всеми горутинами.

func main() {
	go fmt.Println("возможно, не успею напечататься")
	// main завершился — программа вышла, горутина могла не выполниться
}

Значит, нам нужны механизмы синхронизации: дождаться завершения. Их два семейства — каналы и пакет sync.

Каналы

Канал — типизированная труба, по которой горутины передают значения. Это и есть «общение вместо разделения памяти».

ch := make(chan int)   // небуферизованный канал целых
ch <- 42               // отправить (стрелка внутрь канала)
v := <-ch              // получить
close(ch)              // закрыть (делает отправитель, никогда получатель)

Небуферизованные каналы синхронизируют

У небуферизованного канала нет места для хранения. Отправка блокируется, пока кто-то не начнёт получать, и наоборот. Это делает канал точкой рандеву — механизмом синхронизации, а не только передачи данных.

func main() {
	done := make(chan struct{})   // chan struct{} — сигнальный канал, не несёт данных
	go func() {
		fmt.Println("работаю...")
		done <- struct{}{}        // сигналим о завершении
	}()
	<-done                        // main блокируется здесь, пока не придёт сигнал
	fmt.Println("готово")
}

chan struct{} — идиома для «сигнала без данных»: пустая структура занимает 0 байт.

Буферизованные каналы

Буферизованный канал вмещает N значений. Отправка блокируется только когда буфер полон, приём — когда пуст. Буфер полезен, чтобы отправитель не ждал получателя при всплесках, и для ограничения параллелизма (семафор).

ch := make(chan int, 3)  // буфер на 3
ch <- 1; ch <- 2; ch <- 3 // не блокируется
ch <- 4                   // заблокируется — буфер полон

Итерация и закрытие

Закрытие канала — сигнал «данных больше не будет». Получатели узнают об этом через comma-ok или for range:

v, ok := <-ch     // ok == false, если канал закрыт и пуст

for v := range ch { // читает, пока канал не закроют; тогда цикл завершится
	process(v)
}

Правила, нарушение которых карается паникой или дедлоком:

  • Закрывает только отправитель, и только один. Закрыть уже закрытый канал — паника. Записать в закрытый — паника.
  • Читать из закрытого канала можно — вернётся нулевое значение с ok == false.
  • Чтение/запись в nil-канал блокируется навсегда (иногда используется намеренно в select).

select — переключатель каналов

select ждёт сразу несколько канальных операций и выполняет ту, что готова первой. Если готовы несколько — выбирает случайную. Это ядро всех конкурентных паттернов.

select {
case v := <-ch1:
	fmt.Println("из ch1:", v)
case ch2 <- 42:
	fmt.Println("отправили в ch2")
case <-time.After(time.Second):
	fmt.Println("таймаут — за секунду ничего не случилось")
default:
	fmt.Println("ничего не готово прямо сейчас (неблокирующий select)")
}
  • time.After даёт канал, срабатывающий через заданное время — так делают таймауты.
  • default превращает select в неблокирующий: если ничего не готово, идём в default.

context — отмена и дедлайны

Каналы отвечают за передачу данных, а context.Context — за распространение сигнала отмены и дедлайнов сквозь дерево вызовов и горутин. Это стандарт де-факто: любая функция, делающая I/O (запрос к БД, HTTP-вызов), должна принимать ctx context.Context первым аргументом.

Зачем это нужно на первых принципах: пользователь закрыл соединение — не надо продолжать тяжёлый запрос к БД. Истёк таймаут — надо остановить всю цепочку работы, а не только верхний вызов. Контекст — это тот самый провод отмены, протянутый через всю операцию.

// Контекст с таймаутом на всю операцию
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()   // ВСЕГДА вызывайте cancel, иначе утечка ресурсов контекста

result, err := db.QueryContext(ctx, "SELECT ...")
// если за 3 секунды не успели — QueryContext вернёт context.DeadlineExceeded

Как реагировать на отмену внутри своей горутины — слушать ctx.Done():

func worker(ctx context.Context, jobs <-chan Job) error {
	for {
		select {
		case <-ctx.Done():
			return ctx.Err() // context.Canceled или context.DeadlineExceeded
		case job, ok := <-jobs:
			if !ok {
				return nil
			}
			if err := process(ctx, job); err != nil {
				return err
			}
		}
	}
}

Конструкторы контекста:

  • context.Background() — корень, обычно в main.
  • context.WithCancel(parent) — ручная отмена вызовом cancel().
  • context.WithTimeout(parent, d) — отмена по истечении d.
  • context.WithDeadline(parent, t) — отмена в конкретный момент времени.
  • context.WithValue(parent, k, v) — прокинуть значение (request ID, трассировку). Не злоупотребляйте: контекст не помойка для параметров.

Правила гигиены: контекст — первый параметр, не храните его в структурах, всегда вызывайте cancel (лучше через defer), не передавайте nil — используйте context.TODO() как заглушку. Полное введение — Go Concurrency Patterns: Context в официальном блоге.

Пакет sync — когда каналы избыточны

Иногда «общаться через каналы» — из пушки по воробьям. Защитить счётчик или дождаться группы горутин проще примитивами из sync. Мантра: используйте каналы для передачи владения данными, мьютексы — для защиты общего состояния.

WaitGroup — дождаться группу горутин

var wg sync.WaitGroup
for _, url := range urls {
	wg.Add(1)                 // +1 к счётчику ДО запуска горутины
	go func(u string) {
		defer wg.Done()       // -1 при выходе
		fetch(u)
	}(url)
}
wg.Wait()                     // блокируется, пока счётчик не станет 0

Классическая ловушка (до Go 1.22) — захват переменной цикла: раньше go func(){ fetch(url) }() без параметра печатал последний url для всех итераций. В Go 1.22 семантику переменной цикла починили (каждая итерация — своя переменная), но передавать значение параметром — по-прежнему хороший, явный стиль.

Mutex — защита общего состояния

type Counter struct {
	mu sync.Mutex
	n  int
}

func (c *Counter) Inc() {
	c.mu.Lock()
	defer c.mu.Unlock()
	c.n++                     // критическая секция под замком
}

sync.RWMutex разделяет читающие и пишущие блокировки: много читателей одновременно, писатель — эксклюзивно. Полезен при «много чтений, редкие записи». Не забывайте: нулевое значение мьютекса готово к работе, инициализировать не надо. И никогда не копируйте структуру с мьютексом (go vet это ловит).

Once — ровно один раз

var (
	once     sync.Once
	instance *Config
)

func GetConfig() *Config {
	once.Do(func() {          // тело выполнится ровно один раз, даже при гонке горутин
		instance = loadConfig()
	})
	return instance
}

errgroup — WaitGroup, который умеет ошибки

Пакет golang.org/x/sync/errgroup — не стандартная библиотека, но почти обязательная. Он как WaitGroup, но собирает первую ошибку и отменяет общий контекст при сбое любой горутины.

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

func fetchAll(ctx context.Context, urls []string) error {
	g, ctx := errgroup.WithContext(ctx)
	for _, url := range urls {
		url := url
		g.Go(func() error {           // возвращаем ошибку прямо из горутины
			return fetch(ctx, url)
		})
	}
	return g.Wait()  // вернёт первую ненулевую ошибку; ctx отменится при первом сбое
}

g.SetLimit(n) ограничивает число одновременно работающих горутин — встроенный worker pool. Это, пожалуй, самый удобный инструмент для параллельной обработки с обработкой ошибок в проде.

Конкурентные паттерны

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

Запускать миллион горутин на миллион задач можно, но часто нежелательно: они завалят БД или внешний API. Пул воркеров фиксирует число одновременных обработчиков.

func workerPool(ctx context.Context, jobs <-chan Job, results chan<- Result, workers int) {
	var wg sync.WaitGroup
	for i := 0; i < workers; i++ {
		wg.Add(1)
		go func() {
			defer wg.Done()
			for job := range jobs {           // все воркеры читают из одного канала
				select {
				case <-ctx.Done():
					return
				case results <- process(job):
				}
			}
		}()
	}
	wg.Wait()
	close(results)                            // закрываем результаты, когда все воркеры отработали
}

Fan-out / Fan-in

Fan-out — раздать работу из одного канала множеству горутин (это и есть пул выше). Fan-in — слить результаты множества горутин в один канал.

// Fan-in: объединяем несколько входных каналов в один выходной
func fanIn(ctx context.Context, chans ...<-chan int) <-chan int {
	out := make(chan int)
	var wg sync.WaitGroup
	for _, c := range chans {
		wg.Add(1)
		go func(c <-chan int) {
			defer wg.Done()
			for v := range c {
				select {
				case out <- v:
				case <-ctx.Done():
					return
				}
			}
		}(c)
	}
	go func() { wg.Wait(); close(out) }()  // закрыть out, когда все источники иссякнут
	return out
}

Pipeline — конвейер стадий

Данные текут через цепочку стадий, каждая — горутина, соединённая каналами. Каждая стадия читает из входного канала, обрабатывает, пишет в выходной.

// Стадия 1: генератор
func gen(nums ...int) <-chan int {
	out := make(chan int)
	go func() {
		defer close(out)
		for _, n := range nums {
			out <- n
		}
	}()
	return out
}

// Стадия 2: возведение в квадрат
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() {
	for n := range sq(sq(gen(2, 3))) {  // 2->4->16, 3->9->81
		fmt.Println(n)
	}
}

Эталонная статья с этим примером и правильной обработкой отмены — Go Concurrency Patterns: Pipelines and cancellation.

Гонки данных и go test -race

Гонка данных (data race) — это когда две горутины обращаются к одной переменной без синхронизации и хотя бы одна пишет. Результат — неопределённое поведение: испорченные данные, случайные падения, баги, которые «не воспроизводятся».

// ГОНКА: несколько горутин пишут в counter без защиты
counter := 0
for i := 0; i < 1000; i++ {
	go func() { counter++ }()  // counter++ не атомарна: read-modify-write
}

Гонки коварны тем, что код часто «работает» на разработке и падает в проде под нагрузкой. Поэтому в Go есть встроенный детектор гонок. Запускайте тесты и бинарники с флагом -race:

go test -race ./...
go run -race ./cmd/server
go build -race -o server ./cmd/server

Детектор инструментирует доступы к памяти и в рантайме сообщает: какие горутины, к какому адресу, откуда. -race в CI на тестах — обязательная практика для любого конкурентного кода. Замедляет исполнение в 5-10 раз, поэтому в проде его не гоняют, только в тестах. Подробно — Data Race Detector.

Способы починки: защитить sync.Mutex, использовать sync/atomic (atomic.Int64 с Go 1.19), или переструктурировать так, чтобы к переменной обращалась только одна горутина (общение через каналы).

Утечки горутин

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

// УТЕЧКА: горутина навсегда застряла на отправке, потому что получатель ушёл по таймауту
func leak() {
	ch := make(chan int)      // небуферизованный
	go func() {
		ch <- 42              // заблокируется НАВСЕГДА, если никто не прочитает
	}()
	// функция вернулась, читателя нет — горутина висит вечно
}

Правила против утечек:

  • У каждой запущенной горутины должен быть гарантированный путь завершения (закрытие канала, ctx.Done(), таймаут).
  • Передавайте context.Context и слушайте ctx.Done() в долгоживущих горутинах.
  • В select на отправку добавляйте ветку <-ctx.Done(), чтобы не застрять.
  • Диагностика: runtime.NumGoroutine(), профиль pprof (go tool pprof), детектор в тестах — go.uber.org/goleak проверяет, что тест не оставил висящих горутин.

Graceful shutdown

Продакшн-сервис должен уметь корректно останавливаться: перестать принимать новые запросы, дообработать текущие, закрыть соединения с БД — и всё это в отведённое время. Это прямое применение context и сигналов ОС.

func main() {
	srv := &http.Server{Addr: ":8080", Handler: router()}

	// Слушаем сигналы ОС; ctx отменится по SIGINT/SIGTERM
	ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
	defer stop()

	go func() {
		if err := srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
			log.Fatalf("сервер упал: %v", err)
		}
	}()
	log.Println("сервер запущен на :8080")

	<-ctx.Done()  // блокируемся до сигнала завершения
	log.Println("получен сигнал остановки, завершаемся...")

	// Даём 10 секунд на до-обработку текущих запросов
	shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
	defer cancel()
	if err := srv.Shutdown(shutdownCtx); err != nil {
		log.Printf("принудительное завершение: %v", err)
	}
	log.Println("сервер остановлен")
}

srv.Shutdown перестаёт принимать новые соединения и ждёт завершения активных до истечения shutdownCtx. Это разница между «уронили 200 запросов пользователей при деплое» и «ни один не заметил».

Ментальные модели, которые стоит закрепить

  • Каналы — для передачи владения данными между горутинами и синхронизации.
  • Mutex/atomic — для защиты общего состояния, когда передавать нечего.
  • context — сквозной провод отмены и дедлайнов; первый аргумент I/O-функций.
  • errgroup — параллелизм с обработкой ошибок и лимитом; берите его по умолчанию для «сделай N вещей параллельно».
  • Каждая горутина должна иметь известный способ завершиться.
  • -race в CI — не опция, а обязанность.

Лучшие источники по конкурентности

Что дальше

Мы умеем писать корректный конкурентный код. Но как доказать, что он корректен? Переходим к тестированию — включая -race, табличные тесты и моки.

05. Тестирование

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

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

Доска запросов