Back-pressure pattern: explicit flow control between producer and consumer. Demonstrates 4 scenarios — no back-pressure leading to OOM, Reactive Streams request(N) flow, TCP receive window back-pressure, and bounded queue with drop policy for metrics ingest.
Back-pressure нужен в любой системе, где одна часть производит работу быстрее, чем другая может ее обработать. Без явного flow control разница скоростей не исчезает, она превращается в очередь, рост памяти, рост latency, retry storm или потерю данных. На маленькой нагрузке это выглядит как безобидный буфер. На production-пике это становится OOM, зависшими worker threads и cascade failure.
Классическая ситуация: producer генерирует 10 тысяч событий в секунду, consumer обрабатывает тысячу. Каждую секунду появляется 9 тысяч лишних событий. Через минуту это уже сотни тысяч объектов в RAM. Если очередь unbounded, система выглядит стабильной ровно до момента, когда ее убивает memory pressure. Если очередь bounded, но нет политики overflow, система начинает блокировать случайные участки pipeline. Если есть drop, но нет метрик, команда не знает, сколько данных потеряла.
Back-pressure делает перегрузку частью протокола: consumer явно сообщает, сколько готов принять, producer уважает лимит, а система деградирует контролируемо. Это не только Reactive Streams. TCP receive window, HTTP/2 flow control, gRPC streaming, bounded channels в Go, pause/resume в Node.js streams и consumer lag в Kafka — разные проявления одной идеи.
Mental model простой: не producer решает, сколько можно отправить, а downstream сообщает свою емкость. Consumer говорит: request(100), значит upstream может отправить не больше 100 элементов. Когда batch обработан, consumer просит еще. Скорость всей цепочки автоматически становится скоростью самого медленного полезного звена.
Это отличается от простого rate limiting. Rate limit часто задается извне: не больше N запросов в секунду. Back-pressure рождается из фактического состояния downstream: queue depth, свободные connections, TCP window, lag, доступная память, скорость flush в storage. Поэтому back-pressure лучше отражает реальность, но требует, чтобы producer умел замедляться, отменять работу или принимать отказ.
ADR-формулировка: между async-стадиями запрещаем unbounded queue. Для каждого буфера фиксируем capacity, overflow policy, метрики и поведение producer при сигнале перегрузки. Для streaming API выбираем протокол с явным demand или встраиваем demand-сигнал на уровне приложения.
Диаграмма содержит producer, bounded queue, consumer и storage. Producer заявлен как 10K events/s, consumer как 1K events/s. Этого достаточно, чтобы увидеть фундаментальный конфликт: без регулирования backlog растет на 9K events/s. Очередь находится не как абстрактная деталь, а как явный компонент с decision-записью: где возможен request(N), где нужен bounded queue, где допустим drop-oldest, а где нужно блокировать или отправлять в dead-letter.
Связи показывают естественный pipeline: producer enqueue, queue dequeue, consumer persist. В сценарии TCP тот же рисунок переиспользуется как server/client flow: receive window уменьшается, write блокируется, worker thread висит. Это специально полезно: back-pressure не привязан к очередям, он возникает на каждом слое, где есть разница между скоростью отправки и скоростью приема.
Без back-pressure — OOM демонстрирует, почему unbounded queue является отложенной аварией. Пока память есть, система кажется успешной: producer не блокируется, consumer работает, storage принимает записи. Но backlog растет быстрее, чем его можно разобрать. Урок: если размер буфера не ограничен, пределом становится heap, disk или лимит процесса, а не архитектурное решение.
Reactive Streams — request(N) показывает consumer-driven demand. Consumer подписывается и просит 100 элементов. Producer отправляет ровно 100 и ждет следующего запроса. Система не требует угадывать sleep между batches: demand сам течет вверх по pipeline. Урок: для stream processing лучше протокол, где demand является частью контракта.
TCP receive window показывает встроенную back-pressure в сети. Slow client заполняет receive buffer, ACK сообщает все меньший window, потом window становится нулем. Blocking write на сервере может удержать thread и превратить медленных клиентов в DoS. Урок: back-pressure на TCP-уровне не отменяет application-level timeouts, async I/O и лимиты на slow connections.
Bounded queue + drop oldest показывает осознанную потерю данных для metrics ingest. На пике очередь заполняется, старые метрики удаляются, новые принимаются. Это разумно для telemetry, где свежие данные ценнее старых, но опасно для финансов, audit log и ordered domain events. Урок: overflow policy должна соответствовать семантике данных.
ADR: Buffer vs back-pressure. Буфер сглаживает короткие всплески и decouple-ит producer от consumer. Но буфер не создает capacity, он только откладывает проблему. Back-pressure заставляет upstream замедлиться и сохраняет память, но делает зависимость между стадиями видимой. Для bursty workload обычно нужен небольшой bounded buffer плюс явный сигнал перегрузки.
ADR: Drop vs block. Drop подходит для metrics, logs, clickstream sampling и других данных, где потеря части событий лучше, чем остановка сервиса. Block подходит для недопустимой потери: платежи, ledger, audit, инвентарь. Но block переносит проблему upstream и может создать cascade failure. Поэтому block почти всегда должен иметь timeout и понятный отказ.
ADR: Drop oldest vs drop newest. Drop oldest сохраняет актуальное состояние и подходит для мониторинга. Drop newest сохраняет историческую последовательность и может быть полезен, если consumer обязательно должен разобрать старые элементы. Для ordering-sensitive потоков лучше не drop, а durable queue, backoff producer или dead-letter.
ADR: Reactive demand vs durable log. Reactive Streams хороши для online pipeline с живым producer и consumer. Kafka-style durable log лучше, когда consumer может отставать минуты или часы, а данные нельзя потерять. Но Kafka не отменяет back-pressure: consumer lag, partition count, retention, producer acks и broker disk throughput становятся частью capacity-плана.
ADR: Autoscaling vs flow control. Autoscaling добавляет capacity, но с лагом: image pull, warmup, rebalance, cache warm. Back-pressure действует немедленно. Надежная система обычно использует оба механизма: back-pressure для мгновенной защиты, autoscaling для восстановления headroom.
Reactive Streams формализуют demand через Subscription.request(N) и используются в Akka Streams, Project Reactor и RxJava. gRPC наследует HTTP/2 flow control, где есть per-stream и connection-level окна. TCP receive window есть почти в любой сетевой системе, но он часто проявляется неприятно: synchronous server может зависнуть на write к медленному клиенту.
Kafka решает часть проблемы durable-логом: producer пишет в broker, consumer читает в своем темпе, lag становится измеримой величиной. Но если lag растет быстрее retention, данные будут потеряны или consumer никогда не догонит. Node.js streams имеют pause()/resume() и сигнал drain, но их легко обойти, если писать в поток и игнорировать return value. Go channels дают естественный block на full buffered channel, но без timeout это может заблокировать всю цепочку.
Главный anti-pattern — unbounded queue между async-стадиями. Она выглядит удобно в тестах, потому что ничего не блокируется, но в production скрывает overload до аварии. Второй anti-pattern — считать, что увеличение буфера решает проблему. Большой буфер часто лишь увеличивает tail latency и делает отказ более дорогим.
Еще одна ошибка — не иметь метрик. Минимальный набор: queue depth, enqueue rate, dequeue rate, consumer lag, drop count, processing latency, retry count, blocked time и age самого старого элемента. Без age можно иметь небольшую очередь из очень старых сообщений и не заметить нарушение SLA.
Опасна и универсальная drop policy. Нельзя использовать один drop-oldest для telemetry и платежей. Нельзя блокировать forever в request path. Нельзя retry-ить без jitter и budget, потому что retry storm добавит еще больше входящих событий в уже перегруженный consumer.
Также часто забывают, что back-pressure должен доходить до источника. Если consumer замедляет queue, но producer продолжает принимать HTTP-запросы без ограничений, перегрузка просто переместится в API layer или load balancer.
Не стоит строить сложный back-pressure protocol для простой batch-задачи, которую можно запустить с фиксированной concurrency и понятным retry. Не нужно вводить Reactive Streams только ради одного cron job.
Back-pressure не заменяет capacity planning. Если средняя входящая нагрузка постоянно выше средней скорости обработки, flow control лишь стабилизирует отказ: очередь не взорвется, но работа будет отклоняться, блокироваться или копиться в durable log. Нужно увеличивать consumer capacity, менять алгоритм, шардировать, кешировать или снижать входной поток.
Не используйте drop-политики для данных с юридической, финансовой или доменной обязательностью. Там нужны durable storage, idempotency, replay, dead-letter и ручная обработка ошибок.
message-queues — durable buffering, lag и retention.load-shedding — осознанный отказ вместо полного падения.rate-limiting — внешнее ограничение входящего потока.capacity-planning — расчет устойчивой скорости consumer и headroom.circuit-breaker — защита downstream при ошибках и timeout.streaming-pipeline — применение back-pressure в event processing.request(N).