open atlas
↑ К треку
Go с нуля до senior GO · 10 · 02

Пайплайны под отменой: защищённые send, дисциплина закрытия и каноничная утечка горутин

Стадии пайплайна связаны каналами; каждый send обязан идти через select с ctx.Done, иначе горутина течёт, едва потребитель перестал читать. Закрывает канал только отправитель; fan-in закрывает после WaitGroup. Буферы гасят джиттер, а не дисбаланс.

GO Senior ◷ 18 min
Уровень
ОсновыJuniorMiddleSenior

Ночной батч сверки работал 14 месяцев. Потом кто-то добавил шаг валидации в последнюю стадию, и на плохом входе она сделала естественное: залогировала и вернулась. Одна битая запись в 02:13, горутина-потребитель выходит — а три стадии выше по течению этого не замечают, им и не положено. Они продолжают производить, пока каждая не заблокируется на отправке в канал, которую никто никогда не примет. За один прогон так замерзало около 40 000 горутин посреди send: ~160 МБ стеков плюс все распарсенные записи, которые они держали. Батч «завершался» (планировщик видел возврат главной горутины), просто память не возвращалась. Через восемь ночей под начал ловить OOMKill в 03:00, и goroutine-профиль читался как чистосердечное признание: один стек, chan send, счётчик 312 887. Паттерн пайплайнов — из поста в блоге Go 2014 года; утечка — из забытого единственного правила, о котором этот пост на самом деле и написан: отправка, которую нельзя бросить, — это горутина, которая не может выйти.

Стадии и каналы — форма

Пайплайн — это горутины, связанные каналами: каждая стадия принимает значения сверху, преобразует и отправляет вниз. Версия из блога 2014 года использовала канал done на пайплайн; современная форма — context:

func parse(ctx context.Context, in <-chan []byte) <-chan Record {
	out := make(chan Record)
	go func() {
		defer close(out) // стадия владеет out — закрывает только отправитель
		for raw := range in {
			rec, err := decode(raw)
			if err != nil {
				continue // или: сообщить в канал ошибок, см. ниже
			}
			select {
			case out <- rec:
			case <-ctx.Done():
				return // бросаем отправку; defer закроет out
			}
		}
	}()
	return out
}

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

Каноничная утечка: неприкрытый send

Уберите select — стадия станет читаться чище, out <- rec, — и станет неправильной. Отправка в канал блокируется, пока получатель не готов. В момент, когда нижняя стадия перестала читать — упала, вернулась раньше на плохом входе, как в Hook, или просто решила, что ей хватит, — эта отправка блокируется навсегда. Горутина припаркована на chan send, GC её не тронет (припаркованная горутина — корень, пришпиливающий её стек и каждую запись, на которую он ссылается), и нигде не всплывает никакая ошибка. Это самая частая реальная утечка горутин, и правило против неё механическое: каждая отправка в канал пайплайна сидит внутри select с веткой ctx.Done(). Приём через for range безопасен лишь потому, что defer close отправителя гарантирован — если стадия может выйти, не закрыв выход, нижние range текут точно так же.

Викторина

В трёхстадийном пайплайне отправки не защищены (голая отправка в out, без select). Конечный потребитель натыкается на битую запись и возвращается раньше, ничего не отменяя. Каково установившееся состояние?

Дисциплина закрытия и fan-in

Когда вы впервые собираете fan-in, вопрос «кто закрывает объединённый канал?» легко решить неправильно: каждый отправитель кажется кандидатом, но любой из них, закрыв канал раньше времени, вызовет панику у остальных. У правила закрытия нет исключений: канал закрывает только отправитель, потому что отправка в закрытый канал паникует и повторное закрытие паникует — получатель в принципе не может знать, что закрывать безопасно. Стадии с одним отправителем получают это бесплатно через defer close(out). Fan-in — N воркеров, шлющих в один общий канал, — место, где дисциплина обретает форму: закрывать нельзя, пока не закончили все отправители, поэтому их считают:

func merge(ctx context.Context, ins ...<-chan Record) <-chan Record {
	out := make(chan Record)
	var wg sync.WaitGroup
	for _, in := range ins {
		wg.Add(1)
		go func() { // по форвардеру на вход
			defer wg.Done()
			for rec := range in {
				select {
				case out <- rec:
				case <-ctx.Done():
					return
				}
			}
		}()
	}
	go func() { wg.Wait(); close(out) }() // WaitGroup, затем close — в отдельной горутине
	return out
}

Пара wg.Wait(); close(out) живёт в отдельной горутине, чтобы merge вернулся сразу — закрытие случается ровно один раз, ровно после выхода последнего форвардера, чем бы этот выход ни был: «вход дочитан» или «ctx отменён посреди send».

Осушать или бросать — кто кого разблокирует при отмене

Отмена разбирает пайплайн с обоих концов, и стоит знать, на какой механизм вы полагаетесь. Бросание — дефолт выше: каждый отправитель защищает свой send через ctx.Done, и при отмене каждая стадия разблокирует себя сама и возвращается; значения в полёте теряются. Осушение — альтернатива: потребитель продолжает for range до закрытия канала, верхние стадии завершаются естественно и ничего не теряется — ценой обработки (или хотя бы приёма) всего уже произведённого. Осушение без защищённых отправок — способ убедить себя, что защита необязательна; оно держится, лишь пока дисциплинирован каждый потребитель, что и сломал один ранний return из Hook. Надёжная стойка: защищайте каждый send в любом случае, а осушение выбирайте, только когда данные нельзя терять (и тогда сам цикл осушения нуждается в границе — счётчик или дедлайн, — иначе медленный производитель превращает остановку в зависание).

Почему это работает

Почему утечка указывает на отправителя, а не на ушедшего получателя? Потому что в каналах CSP-стиля отправка — это обязательство: отправитель посреди send не умеет ни таймаутиться, ни проверить флаг, ни попасть под сборку мусора — его единственные выходы: пришедший получатель или сработавшая другая ветка select. Получатель, вернувшийся раньше, не совершает преступления, видимого системе типов; он просто отнимает будущее, которого ждал отправитель. Эта асимметрия — весь аргумент за идиому защищённого send: спасти отправителя может только он сам, поэтому запасной выход должен быть вкомпилирован в место отправки.

Буферы: что они покупают, честно

Сделать каналы стадий буферизованными — make(chan Record, 1024) — выглядит как фикс пропускной способности. Будьте точны в том, что это делает. Буфер поглощает джиттер: короткие всплески, паузы GC, на миг замедлившегося потребителя; производители движутся дальше вместо передачи в ногу, и латентность небуферизованного рандеву (~сотни наносекунд на пару send-receive) перестаёт доминировать на мелких элементах. Чего буфер не чинит — дисбаланс: если стадия производит 12 тыс./с, а следующая потребляет 9 тыс./с, буфер на 1024 слота полон через ~340 мс, и пайплайн идёт ровно на 9 тыс./с — как небуферизованный, плюс 1024 записи лишней памяти и устаревания. Установившаяся пропускная способность — это самая медленная стадия, точка. Буферы ещё и оттягивают симптомы утечки: неприкрытый send в большой буфер отказывает, только когда буфер заполнится, нередко через минуты после смерти потребителя — баг из 02:13 всплывает в 02:47 со стеком, который уже не указывает на причину.

Викторина

Стадия производит 12 тыс. записей/с; следующая часами потребляет 9 тыс./с. Инженер вставляет между ними буфер на 1024 слота. Что изменится?

Ошибки и наблюдаемость

Стадии, которая может отказать, нужен путь для ошибок. Два рабочих дизайна. Канал ошибок на стадию: стадия шлёт ошибки в errc, оркестратор делает select по нему — явно, но каждый потребитель обязан осушать errc, иначе вы построили вторую утечку. Стадии под errgroup — обычно лучший дефолт: каждая стадия бежит как g.Go, возвращает свою ошибку, а групповой ctx отменяет защищённые отправки всех остальных; g.Wait отдаёт первый отказ после того, как весь пайплайн размотался. Для видимости инструментируйте сами каналы: периодический gauge len(ch) против cap(ch) по стадиям показывает, где сидит очередь — буфер, прибитый к ёмкости, метит стадию-бутылочное-горлышко, буфер, прибитый к нулю, метит голодание выше по течению, а растущий runtime.NumGoroutine() при ровном трафике — сигнал утечки из прошлого урока. Дешёвые числа, превращающие «батч тормозит» в «стена — стадия 3».

Вспомните перед уходом
  1. 01
    Объясни каноничную утечку пайплайна: механизм, почему молчат и рантайм, и GC, и идиому, которая её предотвращает.
  2. 02
    Сформулируй дисциплину закрытия, как её реализует fan-in, и выбор осушать-или-бросать при отмене.
Итог

Паттерн пайплайна — стадии-горутины, связанные каналами, и в современной форме сквозь каждую стадию продет context. Композирующая структура: стадия создаёт и владеет своим выходным каналом, отправляет в него и гарантирует defer close(out), чтобы низ мог for-range; функция возвращает канал сразу, данные текут после сборки. Правило, сохраняющее жизнь под отказом: каждая отправка — select с веткой ctx.Done, потому что неприкрытый send блокируется навсегда, едва низ перестал читать — горутина паркуется на chan send, GC считает её корнем, пришпиливая стек и записи в полёте, рантайм не видит дедлока, пока живы другие горутины, и 40 000 замёрзших отправителей за ночной прогон — цена этой тишины. Дисциплина закрытия — только отправитель; fan-in считает форвардеров WaitGroup-ом и закрывает в отдельной горутине wait-затем-close. Отмена разбирает либо бросанием — каждый отправитель спасает себя через Done, теряя значения в полёте, — либо осушением, где потребитель принимает до close и ничего не теряется, с границей — иначе зависание. Буферы — поглотители джиттера, не лекарство от дисбаланса: буфер на 1024 слота перед стабильно более медленной стадией заполняется за секунды, и пайплайн оседает на скорости самой медленной стадии, платя лишней памятью и оттянутыми симптомами утечек. Гоняйте стадии под errgroup, чтобы один отказ отменял каждую защищённую отправку, а Wait возвращал первую ошибку; экспортируйте len против cap по каналам и число горутин — и пайплайн сам расскажет, где живут его очередь и его утечки. Теперь, когда в goroutine-профиле вы видите сотни одинаковых стеков chan send, — вы знаете, какая стадия осталась без потребителя, и знаете одну строку, которая позволила бы тем горутинам спасти себя самостоятельно.

Практика

Начни сверху. Задачи идут от простого к сложному: вспомнить факт, применить к случаю, затем senior-уровень. Открой, попробуй, потом открой ответ.

вспомнитьприменитьуглубить0 из 6 завершено

Что-то непонятно?

Задай вопрос по этому уроку. Вопросы анонимны и попадают напрямую автору — урок станет лучше.

Примени это

Примени этот урок в реальном проекте.

хоткеи развернуть
поиск
K
пред. пьеса
k
след. пьеса
j
тиры
t
это меню
?
sources3
expand
  1. 01
  2. 02
  3. 03

Trademarks belong to their respective owners. Editorial reference only.