Семафоры и backpressure: ограничение работы на каждой границе, пока очередь не стала аварией
Ограничивай конкурентность семафором на буферизованном канале или semaphore.Weighted; при заполнении блокируй или отказывай — безразмерная очередь это отложенный OOM и сокрытие задержки. Размер — по закону Литтла (L = λW); 503 + Retry-After лучше расплавления.
У сервиса миниатюр была элегантная развязка: запросы добавляли задачи в очередь в памяти, воркеры потребляли в своём темпе. Флеш-распродажа принесла десятикратный всплеск. Очередь сделала свою работу — поглотила всё. Девятнадцать минут дашборд показывал двухсотки и глубину очереди, пробивающую 2.1 миллиона задач по ~1.5 КБ каждая, пока воркеры перемалывали 300 задач в секунду. Поделите: задача, вставшая в очередь на десятой минуте, была бы обработана через два часа — для клиента, который оттаймаутился через 25 секунд. Сервис жёг CPU на работу, которую никто никогда не получит. На 4.2 ГБ под был убит OOM-киллером, перезапустился с пустой очередью, и клиенты — все в ретраях — заполнили её заново за 40 секунд. Три цикла OOM, прежде чем кто-то ограничил очередь и начал возвращать 503. Однострочник из постмортема стал законом команды: безразмерная очередь не поглощает перегрузку — она её откладывает, с процентами и счётом за память.
Два семафора: канал и Weighted
Прежде чем взять любой примитив конкурентности, вопрос один и тот же: что именно ограничено, и где граница находится относительно работы? Семафор без зависимостей — буферизованный канал, где acquire — это отправка:
sem := make(chan struct{}, 64) // ёмкость = максимум в полёте
func handle(job Job) {
sem <- struct{}{} // acquire: блокируется, когда 64 в полёте
defer func() { <-sem }() // release: освобождаем слот
process(job)
}Размещение — деталь уровня senior: acquire до запуска горутины, в цикле производителя — тогда семафор ограничивает и число горутин, и работу. Acquire первой строкой внутри горутины ограничивает только работу: десятикратный всплеск всё равно порождает горутину на задачу, тысячи их паркуются на acquire, каждая ~2–4 КБ плюс захваченное состояние — замедленная версия Hook. Пределы канального семафора: acquire нельзя бросить (нет ветки ctx — если не обернуть в select, что и стоит делать), и все разрешения равны.
golang.org/x/sync/semaphore.Weighted чинит оба: Acquire(ctx, n) возвращается с ошибкой при отмене ctx — застрявший в очереди может из неё уйти, — а веса позволяют работе заявлять свой истинный размер; классика — память: sem.Acquire(ctx, fileSize) против бюджета в байтах, а не счётчика файлов. TryAcquire — вариант быстрого отказа. Честное замечание о цене: Weighted — это мьютекс и список ожидающих; канальная версия чуть дешевле и идиоматична, когда разрешения однородны и вы сами оборачиваете acquire в select.
Краулим 100 тыс. URL с конкурентностью 64. Вариант A берёт канальный семафор до оператора go; вариант B порождает горутину на URL, которая делает acquire первой строкой. Что реально различается при всплеске?
Backpressure: блокировать или отказывать — никогда не копить безгранично
Когда система на пределе, со следующей единицей работы можно сделать ровно три вещи: заблокировать производителя (дефолт семафора — давление распространяется вверх, источник замедляется), отказать быстро (TryAcquire, 503 — вызывающий узнаёт немедленно и может отступить) или поставить в очередь. Первые два честны. Третье честно, только когда очередь ограничена и её задержка посчитана: очередь откладывает работу, а у отложенной работы есть время ожидания, которое прибавляется прямо к вашему времени ответа. Безразмерная версия совершает два греха разом. Она прячет перегрузку от всех сигналов выше — производители видят мгновенные accept, пока система тонет, — и конвертирует перегрузку в рост памяти, который кончается OOM-ом, теряющим всё в полёте, а не только излишек. Двухчасовая задержка очереди из Hook против 25-секундного клиентского таймаута — соль анекдота: за определённой глубиной каждая задача в очереди уже мертва, а воркеры — крематорий с идеальными метриками пропускной способности.
Где живут лимиты: границы
Ставьте границы там, где работа входит, а не там, где болит. На HTTP-краю лимитер запросов в полёте — дюжина строк:
func limitInFlight(n int64, next http.Handler) http.Handler {
sem := semaphore.NewWeighted(n)
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if !sem.TryAcquire(1) {
w.Header().Set("Retry-After", "1") // скажи клиентам, когда вернуться
http.Error(w, "overloaded", http.StatusServiceUnavailable)
return
}
defer sem.Release(1)
next.ServeHTTP(w, r)
})
}Быстрый 503 с Retry-After лучше расплавления: отклонённый запрос стоит микросекунды, держит латентность ровной для запросов, которые вы приняли, и даёт балансировщикам и клиентам настоящий сигнал. Внутри сервиса размер пулов воркеров определяет ограничивающий ресурс, а не число CPU: пул, зовущий внешний API, ограничен rate-лимитом этого API; пул, пишущий в Postgres, ограничен пулом соединений. И выравнивайте слои: db.SetMaxOpenConns(10) уже и есть семафор — принять 500 одновременных запросов, которым всем нужно соединение, значит передвинуть очередь с края (где отказ дёшев) в пул соединений (где запросы уже съели память, дедлайны и по горутине каждый). Внешний лимит должен быть наименьшим честным из тех, что выдерживают нижние слои.
▸Почему это работает
Почему блокировать производителя бывает лучше, чем отказывать? Потому что блокировка молча композируется сквозь цепочку вызовов, которой вы владеете: стадия пайплайна блокируется на отправке, это тормозит стадию выше, та — читателя у источника; управление потоком без единой строки кода, тот же механизм, что у TCP с его окном приёма. Отказывайте, когда производитель не ваш и замедлить его нельзя (публичные клиенты, партнёрские вебхуки) или когда у работы дедлайн, который ожидание сорвёт. Блокируйте, когда производитель внутренний и эластичный. Единственный неверный ответ — невидимый третий вариант, принять-и-копить без границы: это блокировка, которую никто не чувствует, и отказ, который приходит на два часа позже в виде OOM.
Закон Литтла: размер по арифметике, а не по ощущениям
Инструмент калибровки — одна строка алгебры: L = λW — средняя конкурентность равна частоте прихода, умноженной на время в системе. Работает во все стороны. Проверка спроса: 400 req/s по 100 мс на запрос — это 400 × 0.1 = 40 в полёте в среднем; лимит 32 мал ещё до первого всплеска, ведь максимум при лимите 32 — это 32 / 0.1 = 320 req/s, и недостающие 80 req/s должны копиться (вечно расти) или отклоняться. Перевод очереди: 2.1 млн задач из Hook при 300 задач/с — это W = L/λ = 7 000 с задержки очереди — видимый абсурд рядом с таймаутом 25 с; именно поэтому глубина очереди на дашборде должна делиться на скорость разбора и показываться в секундах, а не штуках. Запас: берите лимит как спрос плюс запас на всплеск (при среднем 40 в полёте разумен лимит 64), затем проверьте, что каждый нижний слой выдерживает подразумеваемую пропускную способность. Арифметика занимает тридцать секунд и завернула бы дизайн из Hook ещё на ревью.
Сервис стабильно получает 400 req/s; каждый запрос занимает воркера на 100 мс. Лимит в полёте — 32. Что говорит закон Литтла?
Лестница: отказ, сброс, деградация
Лимиты решают, когда вы за пределом; лестница решает, чем пожертвовать. Отказ новой работе на краю — 503 плюс Retry-After, самый дешёвый и честный вариант, первая ступень, потому что защищает всё ниже. Сброс по приоритету — чекаут оставить, префетчи и фоновые обновления выбросить; для этого работа должна классифицироваться на входе — потому теги приоритета живут в запросе, а не в воркере. Деградация самой работы — отдать кешированную миниатюру вместо рендера, пропустить блок рекомендаций, вернуть частичный результат с честным заголовком. Большинство команд открывает ступени в обратном порядке, во время инцидентов; на этапе дизайна дешевле. И в то же ревью входит защита от шторма ретраев: отклонённые клиенты обязаны отступать с джиттером, иначе ваши 503 синхронизируют их в тот самый всплеск, что заполнил очередь Hook за 40 секунд. Вместе все три ступени означают: при перегрузке у вас всегда есть ответ дешевле расплавления — но провести ступени нужно до инцидента, а не во время него.
Ваш HTTP-сервис на пределе (все 64 слота воркеров заняты). Приходит новый запрос. Какая стратегия ответа лучше всего подходит для публичного API с клиентским таймаутом 25 с?
- 01Сравни семафор на буферизованном канале и semaphore.Weighted: механика, ловушка размещения acquire и когда какой выигрывает.
- 02Почему именно безразмерная очередь хуже блокировки и отказа, и как закон Литтла это квантифицирует?
У ограниченной конкурентности в Go две родные формы. Семафор на буферизованном канале — acquire это отправка, release — приём в defer — дефолт без зависимостей, с одним правилом размещения, отделяющим его от замаскированной безразмерной очереди: acquire в цикле производителя, до оператора go, чтобы граница покрывала горутины наравне с работой. semaphore.Weighted добавляет отменяемое ожидание (Acquire принимает ctx) и взвешенные разрешения — канонический случай: бюджет байт для тяжёлой по памяти работы; TryAcquire в обоих даёт путь быстрого отказа. Принцип, которому служат инструменты: на пределе блокируй производителя или отказывай работе — блокировка композируется вверх по цепочкам, которыми ты владеешь, как управление потоком в TCP; отказ с 503 и Retry-After — честный ответ производителям, которых не замедлить. Единственный запрещённый вариант — безразмерная очередь: она прячет перегрузку от всех сигналов, конвертируя её в рост памяти — отложенный OOM плюс сокрытие задержки, и работа в очереди истекает по клиентским таймаутам задолго до того, как воркеры до неё доберутся. Ставь лимиты на границах: лимитер в полёте на HTTP-краю, пулы воркеров по ограничивающему ресурсу, а не по числу CPU, и слои, выровненные так, чтобы внешний лимит не превышал того, что выдержит пул БД — сам по себе семафор. Ничего не калибруй интуицией: L = λW превращает частоту и латентность в требуемую конкурентность, а глубину очереди, делённую на скорость разбора, — в секунды отставания для алерта. Над лимитами — лестница: отказ, сброс по приоритету, деградация работы — плюс клиентский backoff с джиттером, чтобы твои отказы не синхронизировались в следующий всплеск. Теперь, когда видите растущую глубину очереди при стабильной латентности, первым делом считайте W = L/λ: если глубина очереди, делённая на скорость разбора, уже превышает клиентский таймаут — очередь это морг, и граница должна переехать на вход.
Практика
Начни сверху. Задачи идут от простого к сложному: вспомнить факт, применить к случаю, затем senior-уровень. Открой, попробуй, потом открой ответ.
Что-то непонятно?
Задай вопрос по этому уроку. Вопросы анонимны и попадают напрямую автору — урок станет лучше.
Примени это
Примени этот урок в реальном проекте.