Backpressure в очереди уведомлений

26 мая 2026 · бэкенд, очереди

очереди надёжность Redis

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, из него читает пул из восьми воркеров; пунктирная стрелка возвращается от воркеров к источникам как возврат кредита, отдельная стрелка вниз показывает уход низкоприоритетных событий в дайджест
Кредит — это разрешение записать одно сообщение. Пока воркер не забрал предыдущее, нового разрешения не появится.

Реализация вышла скучной, что хорошо. Ёмкость 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 3002 048
Задержка доставки, p998 мин 20 с1,4 с
Задержка критичных, p998 мин 20 с0,6 с
Отказов критичным за месяц0
Схлопнуто рутинных событий0%62% в пике

Компромиссы, о которых стоит знать

  • Отказ действительно возможен. Если критичных событий станет больше, чем 2048 за 250 мс, вызывающий код получит ошибку и покажет её человеку. Это неприятно, но это честно и видно в метрике, в отличие от молчаливого опоздания.
  • Дайджест меняет ощущение продукта. Часть людей заметила, что уведомлений стало меньше, и это понравилось не всем: пропало ощущение «живой ленты». Мы оставили в настройках переключатель для тех, кто хочет каждое событие отдельно.
  • Кредиты — не распределённый механизм. Наш счётчик живёт в памяти процесса. При трёх нодах общая ёмкость получается 3 × 2048, а не 2048, и мы сознательно на это идём: координация ради точной границы дороже, чем небольшая неточность границы.
  • Ёмкость подобрана под текущую скорость пула. Если воркеры замедлятся вдвое, верхняя граница задержки станет 74 секунды. Поэтому у нас теперь алерт не на глубину очереди, а на отношение глубины к пропускной способности — это и есть время разбора.

Главный вывод даже не про очереди. Граница должна быть у любого буфера, и лучше выбрать её самому, чем позволить это сделать памяти процесса или терпению пользователя.