Поток событий поступает быстрее, чем потребитель успевает его обрабатывать. Как механизм обратного давления должен изменить интеграцию?
Нужно ввести обратное давление: скорость публикации или выдачи новых сообщений должна учитывать фактическую способность потребителя. Обычно это означает ограниченные буферы, управление кредитами или размером выборки, ограничение параллелизма и масштабирование потребителей по измеряемому отставанию.
Нельзя решать проблему бесконечным накоплением очереди или молчаливым удалением событий. Для бизнес-событий следует сохранять сообщения до обработки либо явно согласовать допустимую потерю, агрегацию или повторную доставку.
Событийные интеграции стали применять, чтобы разделить отправителя и получателя во времени и переживать кратковременные различия в скорости обработки. Очередь или брокер позволяет производителю не ждать завершения работы потребителя, но такая развязка не устраняет ограниченность его ресурсов.
Если производитель стабильно быстрее потребителя, буфер лишь временно скрывает проблему. Поэтому вместе с асинхронной доставкой появился практический механизм управления потоком: система должна сигнализировать, когда потребитель или промежуточное хранилище приближаются к пределу.
Пусть производитель публикует 10 000 событий в секунду, а потребитель обрабатывает только 7 000. Отставание будет расти на 3 000 событий в секунду, пока не закончится место в очереди, не истечёт срок хранения или не увеличится задержка, неприемлемая для бизнеса.
Без обратного давления возможны переполнение памяти, исчерпание диска, рост времени обработки, каскадные сбои downstream-сервисов и потеря сообщений. Увеличение числа экземпляров потребителя помогает только при наличии независимого параллелизма, достаточных ресурсов и возможности безопасно обрабатывать сообщения одновременно.
Потребитель должен подтверждать или освобождать сообщения только после успешной обработки, а объём одновременно полученных сообщений следует ограничивать. Для этого используют ограниченный размер выборки, лимит необработанных сообщений, ограничение числа параллельных задач или механизм кредитов, при котором потребитель сообщает, сколько сообщений способен принять.
Скорость производителя регулируют через ограничение публикации, временную приостановку, квоты или адаптивное снижение темпа. Если производитель не может быть замедлен, брокер должен обеспечивать долговременное хранение, а система — масштабировать потребителей и контролировать отставание по метрикам consumer lag, времени нахождения события в очереди и доле ошибок.
Обратное давление не равно простому отказу от сообщений. Для обязательных бизнес-фактов безопаснее сохранить событие и обработать его позже. Для телеметрии или промежуточных измерений может быть допустима выборка, агрегация или удаление устаревших значений, но это должно быть частью явного контракта.
Нужно отдельно определить поведение при ошибочном сообщении. Бесконечные повторы одного неисправного события блокируют поток, поэтому применяют ограниченное число попыток, задержку между ними и dead-letter queue или отдельное хранилище ошибок. Это не заменяет исправление причины и не должно приводить к незаметной потере бизнес-события.
Порядок событий также ограничивает масштабирование. Если порядок требуется только внутри одного заказа, пользователя или другого ключа, параллелизм можно строить между ключами, сохраняя последовательную обработку внутри каждого ключа. Если порядок глобальный, увеличение числа потребителей может быть невозможно без изменения требования.
Сервис расчёта рекомендаций получает события о действиях пользователей. Во время рекламной кампании поток вырос втрое, а один тяжёлый тип события начал обрабатываться в несколько раз дольше остальных. Вариант с бесконечной очередью временно скрывал перегрузку, но привёл бы к росту задержки и риску заполнения хранилища.
Рассматривались три решения. Отбрасывать старые события было быстро, но нарушало бы качество аналитики и не подходило для событий, влияющих на обязательные расчёты. Увеличить число потребителей было полезно, однако не устраняло медленный тип обработки и упиралось в лимиты внешнего сервиса. Полностью замедлить всех производителей сохраняло систему, но ухудшало обработку независимых потоков.
Выбрали раздельные очереди для тяжёлых и обычных событий, ограничили число одновременных запросов к внешнему сервису, настроили адаптивную публикацию при росте отставания и ввели отдельную очередь для неисправных сообщений. В результате перегрузка одного класса событий перестала блокировать остальные, а задержка стала контролируемой; цена решения — более сложная маршрутизация и необходимость следить за несколькими показателями отставания.
Нет. Это помогает только если обработка масштабируется горизонтально, сообщения можно выполнять параллельно, а узкие места — CPU, база данных, сеть или внешний сервис — выдерживают дополнительную нагрузку. При требовании последовательности по одному ключу масштабирование возможно между ключами, но не внутри одной последовательности.
Само отставание не является причиной автоматически менять схему события, но оно показывает, что контракт должен учитывать стоимость обработки и допустимую задержку. Полезно разделять обязательные факты и данные, которые можно пересчитать, агрегировать или получать по ссылке. Если потребителю постоянно приходится синхронно запрашивать источник для каждого события, поток может создавать дополнительную нагрузку и терять преимущество асинхронной интеграции.
Только если потеря явно разрешена семантикой конкретного потока. Для событий, фиксирующих бизнес-факты, удаление разрушает возможность восстановить состояние и может привести к несогласованности; предпочтительнее долговременное хранение, повторная обработка или отдельная процедура восстановления. Для метрик, где важна тенденция, а не каждое измерение, допустимы агрегация или выборка, но это должно быть задокументировано в контракте.