Потребитель потока не успевает обрабатывать сообщения. Как механизм обратного давления должен изменить поведение системы?
Обратное давление сообщает источнику или промежуточному брокеру, что потребитель отстаёт, и заставляет ограничить скорость публикации либо временно приостановить её. Само по себе оно не гарантирует сохранность сообщений: для этого нужны буфер, долговечное хранение или явная политика сброса данных.
Потоковые системы отделяют производителей данных от потребителей, поэтому их скорости могут различаться. Без механизма согласования быстрый производитель постепенно переполняет память, очереди или временное хранилище.
Обратное давление появилось как общий способ управлять такой несбалансированной обработкой: потребитель влияет на объём данных, который система разрешает передать ему дальше, вместо бесконтрольного накопления нагрузки.
Пусть источник публикует события быстрее, чем потребитель их обрабатывает. Очередь начинает расти, увеличивается задержка, а при исчерпании памяти возможны ошибки, аварийное завершение процесса или потеря ещё не сохранённых сообщений.
Наивное решение — просто увеличивать буфер. Оно лишь откладывает проблему: если средняя скорость производства надолго выше скорости обработки, конечное хранилище всё равно переполнится. Слишком агрессивное ограничение источника, напротив, может повысить задержку всего конвейера.
Потребитель или промежуточный компонент отслеживает состояние обработки: размер очереди, доступную ёмкость буфера, число сообщений в полёте и скорость подтверждения. При приближении к пределу он уменьшает разрешённый объём данных, снижает частоту запросов, ограничивает число параллельных задач или временно приостанавливает выдачу.
Источник должен уметь принять этот сигнал. В распределённой системе обратное давление может распространяться по цепочке от медленного потребителя к брокеру и далее к производителю. Если источник нельзя замедлить, система должна перенаправить поток в долговечное хранилище, увеличить число потребителей или применить заранее определённую деградацию.
Важно различать обратное давление и гарантию доставки. Обратное давление управляет скоростью прохождения данных, а гарантия доставки определяется подтверждениями, повторной обработкой, хранением сообщений и идемпотентностью потребителя. При сбое часть данных всё равно может быть потеряна, если она находилась только в оперативной памяти.
Основной компромисс — баланс между задержкой, пропускной способностью и ресурсами. Ограничение скорости защищает систему от перегрузки, но повышает время ожидания. Увеличение параллелизма ускоряет обработку, однако может перегрузить внешнюю базу или нарушить порядок событий.
Если сообщения нельзя терять, обычно выбирают долговечный буфер с ограниченным периодом хранения, явным контролем отставания и масштабированием потребителей. Если часть данных допустимо отбросить, например промежуточные метрики, можно применять семплирование или объединение сообщений, но это уже политика качества данных, а не универсальное свойство обратного давления.
Сервис телеметрии публиковал показания быстрее, чем аналитический потребитель записывал их в хранилище. Рассматривались три варианта: бесконечно наращивать память, отбрасывать старые сообщения или ограничивать публикацию при росте отставания.
Наращивание памяти было простым, но не решало проблему длительного дисбаланса и создавало риск аварии. Отбрасывание сообщений сохраняло доступность, но делало невозможным восстановление полной истории. Выбрали долговечную очередь с ограничением числа сообщений в обработке и масштабированием потребителей; при достижении предельного отставания источник временно снижал частоту публикации.
В результате память потребителей перестала расти без ограничения, а задержка стала контролируемой. При этом команда отдельно определила срок хранения очереди и правила поведения при его превышении, поскольку обратное давление не заменяет план аварийного восстановления.
Нет. Оно предотвращает неконтролируемое переполнение, регулируя скорость передачи. Потеря данных исключается только при наличии достаточного долговечного буфера, корректных подтверждений, повторной доставки и обработки ситуации, когда поток не успевает потребляться дольше времени хранения.
Нужно отделить источник от потребителя через промежуточное долговечное хранилище, масштабировать обработку или применить согласованное прореживание данных. Увеличение буфера помогает только на ограниченном интервале; при постоянном превышении скорости производства над скоростью потребления оно лишь переносит отказ на более поздний момент.
Параллелизм полезен, если обработка действительно распараллеливается и узкое место находится в вычислениях. Он не поможет, если ограничивающим ресурсом является одна внешняя база, последовательный раздел потока, сетевое соединение или гарантия порядка. Более того, чрезмерное масштабирование может усилить конкуренцию за общий ресурс и увеличить число ошибок.