Пайплайны под отменой: защищённые send, дисциплина закрытия и каноничная утечка горутин
Стадии пайплайна связаны каналами; каждый send обязан идти через select с ctx.Done, иначе горутина течёт, едва потребитель перестал читать. Закрывает канал только отправитель; fan-in закрывает после WaitGroup. Буферы гасят джиттер, а не дисбаланс.
Ночной батч сверки работал 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».
- 01Объясни каноничную утечку пайплайна: механизм, почему молчат и рантайм, и GC, и идиому, которая её предотвращает.
- 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-уровень. Открой, попробуй, потом открой ответ.
Что-то непонятно?
Задай вопрос по этому уроку. Вопросы анонимны и попадают напрямую автору — урок станет лучше.
Примени это
Примени этот урок в реальном проекте.