Backpressure в очереди уведомлений
14 апреля в 09:12 один человек перенёс еженедельный созвон на 180 участников — сразу на год вперёд. Очередь уведомлений приняла 9 360 событий за сорок секунд и разгребала их до обеда.
Симптом: напоминание после встречи
Первым пришёл не алерт, а сообщение в чат поддержки: «мне напомнили о дейли в 10:37, дейли было в 10:00». Дальше стало понятно, что это не единичный случай. Напоминание «за 10 минут» — событие с жёстким смыслом: доставленное с опозданием на полчаса, оно не просто бесполезно, оно вредно, потому что человек перестаёт им доверять.
В графиках всё было видно постфактум. Глубина потока notify в пике — 40 300 записей. Задержка от постановки до отправки, p99 — 8 минут 20 секунд. Пул отправки при этом работал ровно на своей обычной мощности: восемь воркеров, около 55 сообщений в секунду, никаких ошибок. Просто входящий поток был в двадцать раз больше исходящего, а границы у очереди не было.
Почему очередь была неограниченной
Потому что так проще, и полтора года это работало. Обычный фон — 90 событий в минуту, воркеры простаивают, глубина потока колеблется около нуля. Продюсер выглядел так:
// Было: запись никогда не ждёт и никогда не отказывает.
func Enqueue(ctx context.Context, n Notification) error {
return rdb.XAdd(ctx, &redis.XAddArgs{
Stream: "notify",
Values: n.Fields(),
}).Err()
}
Неограниченная очередь — это не «очередь без потерь», это «очередь, которая теряет свойство своевременности вместо того, чтобы честно отказать». Событие никуда не делось, оно доехало — через восемь минут, когда уже никому не нужно. Мы обменяли видимую ошибку на невидимую деградацию, и она оказалась хуже.
Три класса событий
Прежде чем ставить границу, надо было решить, что делать при её достижении. Ответ зависит от события, поэтому мы разложили всё, что шлём, на три класса:
- Критичные — напоминание о встрече, отмена встречи, приглашение на встречу через час. Терять нельзя, задерживать нельзя. Около 6% потока.
- Рутинные — «вам назначили задачу», «изменился срок», «добавлен комментарий». Терять нежелательно, задержать на несколько минут — не страшно. Около 80%.
- Фоновые — «доска обновлена», «участник вышел из встречи». Полностью схлопываемые: важно последнее состояние, а не каждый шаг. Остальные 14%.
Дальше стало понятно, что переезд повторяющейся встречи порождает почти исключительно рутинные и фоновые события, а критичные тонут в них просто по невезению.
Ограниченный буфер и кредиты
Реализация вышла скучной, что хорошо. Ёмкость 2048 выбрана из простого расчёта: пул выгребает 55 сообщений в секунду, значит полный буфер разбирается примерно за 37 секунд — это верхняя граница задержки, которую мы готовы терпеть.
type Sink struct {
slots chan struct{} // кредиты: len(slots) — занятые места
out chan Notification
digest *Digest
}
func NewSink(capacity int) *Sink {
return &Sink{
slots: make(chan struct{}, capacity),
out: make(chan Notification, capacity),
}
}
func (s *Sink) Offer(ctx context.Context, n Notification) Outcome {
select {
case s.slots <- struct{}{}:
s.out <- n
return Accepted
default:
}
// Кредитов нет: буфер полон. Что дальше — решает класс события.
switch n.Class {
case Critical:
// Ждём слот, но ограниченно: лучше отказать наверх,
// чем принять и доставить через восемь минут.
ctx, cancel := context.WithTimeout(ctx, 250*time.Millisecond)
defer cancel()
select {
case s.slots <- struct{}{}:
s.out <- n
return Accepted
case <-ctx.Done():
metrics.CriticalRejected.Inc()
return Rejected
}
default:
return s.digest.Fold(n)
}
}
// Воркер возвращает кредит ровно после того, как забрал сообщение.
func (s *Sink) take() Notification {
n := <-s.out
<-s.slots
return n
}
Схлопывание вместо потери
Отказ рутинному событию — плохой ответ: человек всё-таки хочет знать, что ему назначили задачу. Поэтому такие события не выбрасываются, а складываются в дайджест по ключу. Ключ — это то, что делает два события «одним и тем же»:
func foldKey(n Notification) string {
// Двадцать правок одной задачи за минуту — это одно
// уведомление о задаче, а не двадцать.
return n.UserID + "|" + n.Kind + "|" + n.EntityID
}
func (d *Digest) Fold(n Notification) Outcome {
d.mu.Lock()
defer d.mu.Unlock()
k := foldKey(n)
if prev, ok := d.items[k]; ok {
prev.Count++
prev.Latest = n.Payload // важно последнее состояние
return Folded
}
d.items[k] = &Item{Count: 1, Latest: n.Payload, First: time.Now()}
return Folded
}
Дайджест выгружается раз в 60 секунд или когда набирается 200 ключей — что наступит раньше. Пользователь вместо двадцати сообщений получает одно: «в задаче „Перенос отчётности“ 20 изменений». На той же апрельской истории с переносом созвона это схлопнуло 62% рутинного потока.
Отдельно мы поставили жёсткий предел самому потоку в Redis, чтобы даже при ошибке в нашем коде он не рос бесконечно:
XADD notify MAXLEN ~ 20000 * kind reminder user 8813 entity mtg_51f2
Что получилось
| Метрика | 14 апреля | После |
|---|---|---|
| Пиковая глубина очереди | 40 300 | 2 048 |
| Задержка доставки, p99 | 8 мин 20 с | 1,4 с |
| Задержка критичных, p99 | 8 мин 20 с | 0,6 с |
| Отказов критичным за месяц | — | 0 |
| Схлопнуто рутинных событий | 0% | 62% в пике |
Компромиссы, о которых стоит знать
- Отказ действительно возможен. Если критичных событий станет больше, чем 2048 за 250 мс, вызывающий код получит ошибку и покажет её человеку. Это неприятно, но это честно и видно в метрике, в отличие от молчаливого опоздания.
- Дайджест меняет ощущение продукта. Часть людей заметила, что уведомлений стало меньше, и это понравилось не всем: пропало ощущение «живой ленты». Мы оставили в настройках переключатель для тех, кто хочет каждое событие отдельно.
- Кредиты — не распределённый механизм. Наш счётчик живёт в памяти процесса. При трёх нодах общая ёмкость получается 3 × 2048, а не 2048, и мы сознательно на это идём: координация ради точной границы дороже, чем небольшая неточность границы.
- Ёмкость подобрана под текущую скорость пула. Если воркеры замедлятся вдвое, верхняя граница задержки станет 74 секунды. Поэтому у нас теперь алерт не на глубину очереди, а на отношение глубины к пропускной способности — это и есть время разбора.
Главный вывод даже не про очереди. Граница должна быть у любого буфера, и лучше выбрать её самому, чем позволить это сделать памяти процесса или терпению пользователя.