JavaRush /Курсы /Go SELF /Протокол остановки: cancel-on-error и закрытие каналов

Протокол остановки: cancel-on-error и закрытие каналов

Go SELF
69 уровень , 3 лекция
Открыта

1. Почему остановка — отдельная задача

Когда вы впервые пишете конкурентную схему, кажется, что всё просто: «в одном месте нашли ошибку — сделали return». И тут Go вежливо (как только Go умеет быть вежливым) напоминает: return выходит только из текущей goroutine, а остальные продолжают жить своей жизнью. Проблема не в том, что они живут, а в том, что они могут продолжать ждать чтение/запись в канал — и ваша программа зависнет так тихо, что даже стыдно: ни panic, ни stack trace, просто вечное ожидание.

Типичная картина «вечного ожидания» выглядит так: collector уже перестал читать results, worker пытается сделать results <- r и блокируется навсегда; координатор ждёт wg.Wait(), но worker не завершится, потому что он заблокирован на отправке; а main стоит и смотрит на это философски. В этот момент вы начинаете понимать, почему конкурентность — это не «быстрее», а «сложнее, но управляемо».

Чтобы остановка была управляемой, нужен протокол: договорённость, кто кому и как сообщает «всё, заканчиваем».

Инструменты остановки: close и context.Cancel

Сейчас у нас есть два очень разных «сигнала остановки», и важно не путать их роли. Первый — это закрытие канала (close(ch)), которое означает: «значений больше не будет». Второй — это отмена контекста (cancel()), которая означает: «операцию надо прекратить как можно скорее».

У context ключевая механика в том, что ctx.Done() возвращает канал-сигнал, который закрывается при отмене; и именно поэтому его удобно слушать в select. Это не магия, а прямой контракт пакета context: Done — это канал отмены, а derived contexts (WithCancel, WithTimeout и т.п.) дают дерево, где отмена родителя отменяет детей.

А close канала — это не «убить всех», а именно «больше отправлять нельзя». Закрытый канал — не мусорка, а табличка “магазин закрыт”. И, как в жизни, если вы пытаетесь «закрыть магазин второй раз», будет скандал (panic: close of closed channel).

2. Владение каналами: кто имеет право делать close

Самая частая причина паники в рабочих пулах — это борьба за право закрывать канал. Прямо как очередь к одному микрофону: все хотят объявить «я закончил», и в итоге микрофон падает.

Запоминаем правило, которое экономит часы жизни: канал закрывает тот, кто в него пишет. Но в worker pool есть нюанс: в results пишет много воркеров, и вот тут возникает «коллективная ответственность», которая в Go обычно заканчивается паникой. Поэтому results закрывает не воркер, а координатор — единственный, кто знает, когда все воркеры точно завершились.

Удобно зафиксировать это маленькой таблицей (и потом держать её в голове, как ПДД):

Канал Кто пишет Кто закрывает Почему именно он
jobs
producer producer он единственный источник задач
results
workers координатор воркеров много, закрывать должен один после wg.Wait()

Если вам хочется закрыть results из воркера, остановитесь и спросите себя: «а остальные воркеры уже точно не попытаются туда написать?» Если ответ не “железобетонно да” — значит, закрывать нельзя.

4. Cancel-on-error: гарантии и места зависаний

Cancel-on-error — это не «при ошибке всё мгновенно исчезает». Это реалистичнее: «при первой серьёзной ошибке мы даём общий сигнал остановки, а все goroutine обязуются иметь путь выхода и не зависать на каналах».

Тут важны две гарантии.

Первая гарантия: producer не должен зависнуть на отправке задач, если воркеры уже начали сворачиваться. Для этого producer отправляет в jobs через select с ctx.Done().

Вторая гарантия: worker не должен зависнуть на отправке результата, если collector уже «устал и ушёл». Для этого worker отправляет в results тоже через select с ctx.Done().

Именно поэтому ctx.Done() так любят: он превращает потенциально вечное ожидание в «ожидание, но с аварийным выходом». Это и есть та самая практическая ценность контекста, ради которой он вообще существует.

Мини-схема: роли и сигналы остановки

Перед кодом полезно один раз увидеть всю систему как схему. Это снижает вероятность, что вы «чините» зависание в одном месте и случайно создаёте его в другом (конкурентность умеет мстить изощрённо).

flowchart LR
    P[producer] -- jobs --> W1[worker 1]
    P -- jobs --> W2[worker 2]
    P -- jobs --> W3[worker 3]

    W1 -- results --> C[collector]
    W2 -- results --> C
    W3 -- results --> C

    K[координатор] -->|wg.Wait| W1
    K -->|wg.Wait| W2
    K -->|wg.Wait| W3
    K -->|"close(results)"| C

    X["cancel()"] -->|"закрывается ctx.Done()"| P
    X -->|"ctx.Done()"| W1
    X -->|"ctx.Done()"| W2
    X -->|"ctx.Done()"| W3

Слабые места здесь ровно два: стрелки jobs и results. Если где-то перестали читать — кто-то перестанет писать и зависнет. Cancel-on-error нужен, чтобы вместо «зависнуть» было «увидеть Done и корректно выйти».

5. Пример: параллельная обработка задач

Чтобы примеры не были «абстрактными воркерами, которые умножают на 2», давайте представим реальную полезную операцию в нашем приложении с задачами: у нас есть пачка задач, и мы хотим нормализовать заголовки перед сохранением/выводом. Например, запрещаем пустые заголовки и убираем лишние пробелы.

Сделаем функцию, которая принимает []Task, параллельно обрабатывает, и при первой ошибке останавливает остальных.

Начнём с моделей (коротко и скучно — но это та скука, которая спасает от хаоса):

package main

import "strings"

type Task struct {
	ID    int
	Title string
}

func normalizeTitle(s string) (string, error) {
	t := strings.TrimSpace(s)
	if t == "" {
		return "", ErrEmptyTitle
	}
	return t, nil
}

Ошибка ErrEmptyTitle пусть будет sentinel, чтобы её можно было узнавать через errors.Is (мы это уже обсуждали в модуле про ошибки, и это реально полезно).

package main

import "errors"

var ErrEmptyTitle = errors.New("empty title")

6. Реализация cancel-on-error в коде

Каркас: context.WithCancel и дисциплина defer cancel()

Теперь — основа протокола. Мы создаём дочерний контекст и cancel, и сразу ставим defer cancel(). Это не «суеверие», а дисциплина: даже при успехе мы корректно освобождаем ресурсы, связанные с контекстом (таймеры/ссылки и т.п.).

package main

import (
	"context"
)

func processBatch(parent context.Context, tasks []Task) ([]Task, error) {
	ctx, cancel := context.WithCancel(parent)
	defer cancel()

	_ = ctx
	_ = tasks

	return nil, nil
}

Да, пока это заготовка. Но важный момент уже есть: у нас есть единый сигнал остановки для всей группы goroutine.

Типы данных: job и result

Теперь введём типы для очереди и результата. Это слегка увеличивает количество букв, зато резко уменьшает количество ошибок «а что за int в этом канале?».

package main

type job struct {
	Idx  int
	Task Task
}

type result struct {
	Idx  int
	Task Task
	Err  error
}

Дальше мы заведём два канала: jobs и results. И сразу договоримся: jobs закрывает producer (в нашем случае отдельная goroutine), а results закрывает координатор после wg.Wait().

Worker: уважение к ctx.Done() и два select

Самая «скользкая» часть cancel-on-error — воркер. Он может зависнуть в двух местах: при чтении задач и при отправке результата. Значит, и там, и там должен быть select с ctx.Done().

package main

import (
	"context"
	"fmt"
	"sync"
)

func worker(ctx context.Context, jobs <-chan job, results chan<- result, wg *sync.WaitGroup) {
	defer wg.Done()

	for {
		select {
		case <-ctx.Done():
			return
		case j, ok := <-jobs:
			if !ok {
				return
			}

			title, err := normalizeTitle(j.Task.Title)
			if err != nil {
				err = fmt.Errorf("task id=%d: %w", j.Task.ID, err)
			}

			r := result{Idx: j.Idx, Task: Task{ID: j.Task.ID, Title: title}, Err: err}

			select {
			case <-ctx.Done():
				return
			case results <- r:
			}
		}
	}
}

Обратите внимание на wrapping ошибки: fmt.Errorf("...: %w", err) делает ошибку частью цепочки причин, и её можно потом проверять через errors.Is. Это как подпись на коробке: «сломалось при обработке задачи 17», но при этом внутри всё ещё видно, что причина — ErrEmptyTitle.

Producer: отправка задач через select и defer close(jobs)

Producer обязан закрыть jobs при любом исходе: и при успехе, и при отмене. Поэтому defer close(jobs) — это прям часть протокола.

package main

import "context"

func produce(ctx context.Context, jobs chan<- job, tasks []Task) {
	defer close(jobs)

	for i, t := range tasks {
		select {
		case <-ctx.Done():
			return
		case jobs <- job{Idx: i, Task: t}:
		}
	}
}

Здесь важный психологический момент: producer не «должен любой ценой отправить все задачи». Если система отменена — он прекращает отправку. Это и есть смысл cancel-on-error: мы экономим время и ресурсы.

Координатор: wg.Wait() и закрытие results один раз

Координатор — это «дежурный по выключению света». Он ничего не считает и не валидирует, его задача скучная: дождаться wg.Wait() и закрыть results. Скучные задачи — самые важные, потому что именно на них держится отсутствие паник.

package main

import "sync"

func closeResultsWhenDone(results chan<- result, wg *sync.WaitGroup) {
	wg.Wait()
	close(results)
}

Да, это крошечная функция. Но она делает ваш протокол закрытия явным и проверяемым глазами.

Collector: первая ошибка, cancel(), и аккуратное завершение

Collector читает результаты и решает, что делать. В cancel-on-error стратегии логика обычно такая: нашли первую ошибку — запомнили её, вызвали cancel(), а дальше продолжаем читать results до закрытия, чтобы корректно завершить схему.

Почему продолжаем читать? Потому что даже при наличии select в воркере часть результатов может «прилететь» уже после отмены (воркер мог успеть отправить), и мы хотим аккуратно завершить цикл range results. Это не «обязательно всегда», но это очень стабильная привычка.

package main

import (
	"context"
	"sync"
)

func processBatch(parent context.Context, tasks []Task) ([]Task, error) {
	ctx, cancel := context.WithCancel(parent)
	defer cancel()

	jobs := make(chan job)
	results := make(chan result)

	var wg sync.WaitGroup
	workers := 3

	for i := 0; i < workers; i++ {
		wg.Add(1)
		go worker(ctx, jobs, results, &wg)
	}

	go produce(ctx, jobs, tasks)
	go closeResultsWhenDone(results, &wg)

	out := make([]Task, len(tasks))
	var firstErr error

	for r := range results {
		if r.Err != nil && firstErr == nil {
			firstErr = r.Err
			cancel()
		}
		if r.Err == nil {
			out[r.Idx] = r.Task
		}
	}

	if firstErr != nil {
		return nil, firstErr
	}
	return out, nil
}

Здесь есть важная инженерная деталь: мы не «выходим из range results сразу после cancel». Мы ждём, пока results закроется координатором, то есть пока все воркеры завершатся. Так мы не оставляем «висящие» goroutine.

7. Где чаще всего ломают протокол

В cancel-on-error обычно ломают одно из двух мест.

Первое — воркер читает for range jobs и отправляет results <- r без select. В мирной обстановке это работает. Но при ранней остановке collector перестаёт читать — и воркер зависает на отправке. Это зависание «передаётся» вверх: wg.Wait() не заканчивается, results не закрывается, main ждёт вечность.

Второе — producer не слушает ctx.Done() и продолжает пытаться отправить задачи в jobs. Но воркеры уже выходят, и очередь может перестать читаться. Producer зависает на jobs <- .... А дальше всё по той же грустной цепочке.

Cancel-on-error не работает «сам по себе». Он работает только если каждая потенциально блокирующая операция (send/receive) имеет аварийный выход через ctx.Done().

Мысленная проверка: кто кому что должен

Чтобы протокол был не набором трюков, а понятной системой, полезно держать в голове три коротких обязательства, как в договоре аренды (только здесь арендатор — goroutine).

  • Producer обязуется закрыть jobs и перестать отправлять задачи при отмене контекста.
  • Worker обязуется иметь путь выхода и не зависать ни на чтении задач, ни на отправке результатов.
  • Координатор обязуется закрыть results ровно один раз и только после завершения всех воркеров.

Collector обязуется не делать ничего, что оставляет систему «наполовину живой». Если он инициирует остановку — он либо продолжает читать до закрытия results, либо гарантирует, что воркеры не могут зависнуть на отправке (а это как раз достигается select при отправке).

8. Типичные ошибки при cancel-on-error и закрытии каналов

Ошибка №1: закрывать results из воркера.
Это почти всегда приводит к panic: close of closed channel, потому что воркеров несколько, и «первый, кто успел» закроет канал, а следующий попытается закрыть ещё раз или отправить значение в уже закрытый канал. Закрывать results должен координатор после wg.Wait().

Ошибка №2: делать отправку результата как results <- r без select с ctx.Done().
В обычном сценарии это выглядит нормально, но при ранней остановке collector может перестать читать, и воркер зависнет на отправке. Правильный паттерн — select { case results <- r: case <-ctx.Done(): return }, чтобы у воркера был аварийный выход.

Ошибка №3: читать задачи через for range jobs, а отмену контекста не учитывать.
range по каналу завершится только когда канал закрыт. Если producer по какой-то причине завис или не закрыл jobs, воркер будет ждать бесконечно. Для cancel-on-error лучше явный цикл с select, где есть ветка case <-ctx.Done(): return.

Ошибка №4: producer пишет в jobs без select и не реагирует на отмену.
Если воркеры начали сворачиваться, producer может повиснуть на jobs <- .... В итоге вы получили «cancel() вроде вызвали, но программа всё равно зависла». Producer должен отправлять задачи через select и иметь ветку case <-ctx.Done(): return.

Ошибка №5: отменили контекст, но забыли про дисциплину defer cancel().
Если вы создаёте WithCancel, но не вызываете cancel() вообще, то в коде начинают появляться странные эффекты: кто-то ждёт сигнал, который никогда не придёт, или держатся ресурсы дольше, чем нужно. Простое правило “создал cancel — сразу defer cancel()” делает поведение предсказуемым.

1
Задача
Go SELF, 69 уровень, 3 лекция
Недоступна
Генератор суммы
Генератор суммы
1
Задача
Go SELF, 69 уровень, 3 лекция
Недоступна
Три работника
Три работника
1
Задача
Go SELF, 69 уровень, 3 лекция
Недоступна
Деление без хвостов
Деление без хвостов
Комментарии
ЧТОБЫ ПОСМОТРЕТЬ ВСЕ КОММЕНТАРИИ ИЛИ ОСТАВИТЬ КОММЕНТАРИЙ,
ПЕРЕЙДИТЕ В ПОЛНУЮ ВЕРСИЮ