Как я искал backpressure и реальный bottleneck в real-time pipeline

Иногда система заканчивает тест с нулём потерянных данных — и это всё равно провал.

С такой ситуацией я столкнулся при работе над одним real-time telemetry pipeline. Система получала большой поток событий через множество WebSocket-соединений, сохраняла raw data, разбирала сообщения, обновляла внутреннее состояние и передавала данные дальше.

Один из длительных тестов выглядел почти идеально:

  • все 14 558 интересовавших меня событий были доставлены;
  • потерь — 0;
  • writer явно не был перегружен;
  • файлы корректно финализировались.

Но run завершился ошибкой:

raw_ingest_queue_full

То есть данные пока ещё не потерялись, а система уже сообщала: эту нагрузку я устойчиво не выдерживаю.

На разбор проблемы ушло несколько итераций. Самым полезным результатом оказался даже не конкретный performance fix, а способ локализовать backpressure, не пытаясь оптимизировать систему вслепую.

Что за система

Если сильно упростить, поток выглядел так:

WebSocket connections
        ↓
raw ingest queues
        ↓
parsing / normalization
        ↓
state processing
        ↓
observer
        ↓
writer

В реальной версии было много источников и типов сообщений, десятки соединений, bounded queues, telemetry, compression и state updates.

В одном из ранних широких тестов система открывала около 50 WebSocket connections и обслуживала 1676 subscriptions, почти полностью занимая четыре доступных CPU core.

Первый такой run закончился проблемой на shutdown: pipeline не успел корректно завершить drain, остались partial artifacts.

После исправления shutdown/finalization следующий запуск завершался чисто:

  • partial files не осталось;
  • write failures — 0;
  • dropped events — 0.

Но примерно через 21 секунду система снова остановилась:

raw_ingest_queue_full

Теперь проблема была уже не в shutdown. Pipeline просто не успевал перерабатывать входящий поток.

Почему просто увеличить очередь — плохая идея

Когда queue заполнена, первая мысль очевидна:

queue_size *= 10

Иногда большая очередь действительно помогает пережить короткий burst. Но если:

producer_rate > consumer_rate

она только откладывает failure.

Для real-time pipeline это особенно опасно: данные могут физически не теряться, но становиться всё старше.

data loss = 0
latency → ∞

Формально данные целы. Практически система уже перестаёт быть real-time.

Поэтому я решил сначала не менять capacity, а понять, кто именно не успевает.

Где видна проблема — ещё не значит, что там её причина

raw_ingest_queue_full показывает место, где проблема стала видимой, но не обязательно root cause.

Причиной мог быть parser, normalization, state processing, observer, writer — или конкуренция всех этих стадий за CPU.

Следующим шагом стала не optimization, а instrumentation.

Для разных стадий я начал собирать:

  • queue capacity;
  • maximum backlog;
  • peak utilization;
  • producer и consumer rates;
  • produced/consumed event counts;
  • per-connection telemetry.

Из этого получилось простое правило:

Если не знаете, на какой стадии растёт backlog, ещё рано оптимизировать.

Первая гипотеза: не успевает downstream

Подозрение естественно падало на observer и writer. Writer выполняет I/O, поэтому казалось логичным, что именно запись тормозит pipeline.

После нескольких изменений система смогла работать значительно дольше. Следующий важный run продолжался около 17 минут.

За это время было доставлено:

14 558 events

Потеря:

0

Writer queue при этом достигала только:

8.86%

Observer тоже имел большой запас.

Но run снова закончился:

raw_ingest_queue_full

А некоторые raw queues достигли:

100%

Это изменило направление расследования.

Проблема находилась уже не в downstream. После снятия одного bottleneck стал виден следующий — раньше, на raw/parser path.

Новый bottleneck после fix — это не обязательно неудача

Представим:

A → B → C → D

Если D обрабатывает 1 000 событий в секунду, а остальные стадии — 10 000, bottleneck будет D.

После оптимизации D может выясниться, что:

B = 4 000 events/sec

B не стал медленнее. Просто раньше его ограничение было скрыто.

Именно это происходило здесь: после разгрузки observer/writer граница нагрузки переместилась к raw ingestion и parser.

Почему zero data loss всё ещё означал FAIL

Если:

delivered = 14 558
lost      = 0

почему run считается неуспешным?

Потому что zero loss описывает только то, что уже произошло.

Заполненная bounded queue показывает, что произойдёт дальше, если нагрузка сохранится.

Поэтому capacity test должен проверять не только:

lost_events == 0

но и:

queue utilization
producer/consumer balance
freshness
bounded latency

В этом случае 100% raw queue utilization было достаточным основанием остановиться fail-closed даже при нулевых потерях.

Следующая цель — parser path

Per-connection telemetry показывала:

observer queue ≈ свободна
writer queue   ≈ свободна
raw queue      = 100%

Значит, дальше нужно было смотреть между raw queue и следующей стадией.

Среди кандидатов оказались:

  • parsing;
  • lock acquisition;
  • telemetry accounting;
  • частые queue.qsize();
  • обновление статистики для каждого события;
  • слишком маленькие processing batches.

По отдельности каждая такая операция может быть дешёвой. Но внутри hot loop и при сотнях тысяч вызовов этот overhead начинает иметь значение.

Вместо большой переделки — маленький эксперимент

Мне не хотелось сразу переносить pipeline на multiprocessing или разносить стадии по отдельным process. Пока root cause уточняется, это слишком большой change.

Поэтому следующий шаг был минимальным:

  • уменьшить parser-side accounting overhead;
  • часть telemetry обновлять batch-wise;
  • убрать лишнюю работу с hot path;
  • не увеличивать queue capacity.

После изменения прошли локальные тесты, а затем был запущен 180-секундный smoke на реальном потоке данных.

Результат:

early termination: none

raw queue peak:      14.8%
observer queue:       4.87%
writer queue:         9.74%

Producer и consumer state-event counts совпали:

317 329
317 329

По интересовавшему потоку:

delivered: 2 537
lost:          0

До изменения raw queue доходила до 100%.

После — только до 14.8%.

Это было гораздо важнее очередного lost=0: появился capacity headroom.

Почему headroom важнее красивого zero-loss

Сравним два run:

loss                0
queue utilization 100%

и:

loss                0
queue utilization  15%

В обоих случаях потерь нет. Но во втором система имеет заметный запас по capacity, а в первом уже работает на границе saturation.

После изменения raw queue peak снизился со 100% до 14.8%, при этом zero-loss сохранился.

Именно это было главным результатом: не просто отсутствие потерь, а появление capacity headroom.

Не каждый большой stress test хорошо локализует проблему

Ранний тест одновременно включал несколько источников, разные типы потоков, 1676 subscriptions, 50 connections, raw capture, parsing, state processing, telemetry и compression.

Такой run полезен как stress test, но плохо отвечает на вопрос, где именно находится bottleneck.

Для локализации полезнее двигаться ступенчато:

raw only
↓
raw + parsing
↓
raw + parsing + state
↓
full pipeline

Тогда вместо одного большого FAIL постепенно появляется карта capacity системы.

Что оказалось самым полезным

В ходе расследования проблема несколько раз меняла вид:

shutdown
↓
предполагаемый writer pressure
↓
downstream оказался свободным
↓
raw/parser saturation

Каждый следующий run должен был не просто дать PASS или FAIL, а добавить новую информацию.

Рабочая последовательность в итоге стала такой:

failure
↓
instrument
↓
локализовать bottleneck
↓
минимальное изменение
↓
unit tests
↓
короткий реальный smoke
↓
сравнить capacity metrics
↓
только потом длинный run

Это полезнее, чем увеличивать queue и снова запускать тест, не понимая причины.

Несколько правил, которые я теперь использую

Queue full — это симптом. Нужно искать producer и consumer этой очереди.

Zero loss недостаточно. Для real-time pipeline важны также freshness, bounded latency, bounded backlog и capacity headroom.

Большая queue — не throughput. Она может поглотить burst, но не исправляет устойчивое producer > consumer.

Instrumentation часто важнее optimisation. Если telemetry показала, что вы собирались оптимизировать не тот компонент, она уже окупилась.

Новый bottleneck после fix — не обязательно плохая новость. Часто это означает, что предыдущий действительно устранён.

Один PASS-критерий ничего не гарантирует.

В данном случае:

lost_events == 0

было истинно.

Но:

sustainable_capacity == true

— нет.

Итог

Zero-loss сам по себе не означает, что real-time pipeline устойчив.

В моём случае главный сигнал был не lost=0, а то, что после изменения raw queue peak снизился со 100% до 14.8%.

То есть важен не только факт отсутствия потерь, но и запас до saturation.