ADG Оставить заявку
Блог Интеграции 5 мин чтения

Kafka Streams как шина событий: уходим от REST-вызовов между микросервисами

Внедряем Kafka Streams 1.1 как шину событий в микросервисной архитектуре клиента: event-driven вместо синхронных REST-цепочек снижает связанность и убирает каскадные сбои.

Контекст момента

Apache Kafka 1.1 - Kafka Streams для stream-processing в микросервисной архитектуре как замена синхронных REST-вызовов между сервисами

Клиент пришёл с классической проблемой для микросервисной архитектуры: дюжина сервисов, между которыми выстроилась цепочка синхронных REST-вызовов. Сервис A вызывает B, B вызывает C и D, D вызывает E. Пока всё живое - работает. Как только что-то одно падает или тормозит, цепочка складывается как карточный домик. Таймауты копятся, connection pool исчерпывается, пользователь получает 500 с задержкой в несколько секунд.

Мы предложили переосмыслить взаимодействие: вместо того чтобы сервисы вызывали друг друга - пусть публикуют события и реагируют на события других. Apache Kafka 1.1 с Kafka Streams как раз закрывает этот кейс - и без необходимости тащить отдельный stream-processing кластер вроде Flink или Spark Streaming.

Как это устроено было

Архитектура до: пользователь оформляет заказ, запрос идёт в order-service, который синхронно вызывает inventory-service (проверить остаток), потом pricing-service (пересчитать цену с учётом акций), потом notification-service (отправить подтверждение). Каждый вызов - HTTP, с таймаутом в несколько секунд.

Проблема не только в падениях. Ещё большая проблема - связанность. order-service должен знать адреса всех downstream-сервисов, понимать их API, обрабатывать их ошибки. Добавить новый сервис - нужно менять order-service. Убрать - тоже. Это не микросервисы, это монолит, порезанный на кусочки.

Что такое Kafka Streams и почему не просто Kafka

Kafka как брокер существует давно. Kafka Streams - это библиотека для JVM, встроенная прямо в приложение. Никакого отдельного кластера для обработки потоков - логика обработки живёт в самом сервисе, который подписывается на топики и публикует результаты.

В Kafka 1.1, вышедшей в конце марта, ключевые улучшения касаются именно Streams API: переработан механизм rebalancing, появились improvements в KTable и GlobalKTable. Для нашего кейса это важно - мы активно используем KTable для хранения состояния (например, текущих остатков по товарам).

Основная модель:

StreamsBuilder builder = new StreamsBuilder();

KStream<String, OrderEvent> orders =
    builder.stream("order-events");

KTable<String, InventoryState> inventory =
    builder.table("inventory-state");

// джойним поток заказов с состоянием склада
orders.join(inventory,
    (order, stock) -> enrichWithStock(order, stock))
    .to("order-enriched");

Приложение само является частью Kafka-кластера как consumer group. Масштабируется добавлением инстансов - Kafka сам перераспределяет партиции.

Как переложили архитектуру

Вместо цепочки вызовов - топики событий:

  • order-created - order-service публикует событие при создании заказа
  • inventory-reserved - inventory-service читает order-created, резервирует товар, публикует результат
  • order-priced - pricing-service читает order-created, считает цену, публикует
  • order-confirmed - отдельный агрегирующий сервис дожидается оба события и публикует финальный статус
  • notification-send - notification-service читает order-confirmed и шлёт письмо
flowchart LR
    OS["order-service"] -->|order-created| K[("Kafka")]
    K -->|order-created| IS["inventory-service"]
    K -->|order-created| PS["pricing-service"]
    IS -->|inventory-reserved| K
    PS -->|order-priced| K
    K -->|оба события| AG["aggregator-service"]
    AG -->|order-confirmed| K
    K -->|order-confirmed| NS["notification-service"]

order-service больше не знает что существует inventory-service. Он публикует событие и забывает. inventory-service может лечь - заказы всё равно принимаются, сообщения накапливаются в топике, и когда сервис поднимется - он их обработает. Это фундаментальное отличие от синхронной цепочки.

Что вылезло на практике

Первое - Eventually consistent везде. С синхронным REST у вас ответ либо есть, либо ошибка. С event-driven - вы приняли событие, но подтверждения нет, оно придёт позже. Это требует изменений в UX: пользователь видит «заказ принят, подтверждение на почте» вместо мгновенного результата. Клиент к этому был готов, но надо проговаривать заранее.

Второе - дублирование событий. Kafka гарантирует at-least-once delivery. Если inventory-service упал после обработки, но до коммита оффсета - то же сообщение придёт снова. Нужна идемпотентность на стороне потребителя. У нас inventory-service хранит processed event IDs в базе и пропускает повторы. С Kafka Streams 1.1 есть exactly-once semantics через EXACTLY_ONCE в настройках - но это требует Kafka 0.11+, что у клиента есть, и добавляет overhead. На данном этапе остановились на idempotent consumers.

Третье - отладка усложнилась. Раньше проблему можно было отследить по цепочке HTTP-запросов в логах. Теперь событие проходит через несколько топиков и сервисов, и без distributed tracing понять где что зависло - квест. Пришлось добавить correlation ID в каждое событие и пробрасывать его через всю цепочку. Это не автоматически, это надо думать.

Четвёртое - схемы событий требуют версионирования. Когда order-service меняет структуру события - все потребители должны уметь читать старый формат. Мы ввели Avro-схемы и Schema Registry (Confluent Platform). Это ещё один компонент в инфраструктуре, зато контракт между сервисами явный и версионированный.

Kafka Streams 1.1: что нового потрогали

Из конкретных вещей Kafka 1.1 которые оказались полезны в работе: улучшенный rebalancing при рестарте инстансов - при прежних версиях перераспределение партиций могло подвесить обработку на десятки секунд, теперь быстрее. GlobalKTable с автоматической синхронизацией на все инстансы - использовали для хранения справочников (курсы валют, категории товаров), которые нужны во всех партициях.

Из того что пока не устраивает: мониторинг Kafka Streams через JMX по умолчанию - специфично. Нужно либо выгребать метрики самому, либо ставить что-то поверх. Prometheus JMX Exporter с правильным конфигом решает вопрос, но это ещё одна настройка к инфраструктуре.

Где стоим

Пайплайн для основного flow (создание заказа) запущен. На нём несколько тысяч событий в день - для Kafka это смешная нагрузка, но для команды клиента это первый production опыт с event-driven. Пока смотрим как операционная команда справляется с новой моделью отладки и как ведут себя сервисы при нетипичных ситуациях (дублирование, задержка топиков).

Каскадных сбоев типа «упал один - упала цепочка» за три недели не было. Это само по себе результат. Обработка интеграционных задач такого класса - не серебряная пуля, но для систем где сервисов больше пяти и связи нетривиальные, event-driven снижает связанность принципиально, а не косметически.

Оставшиеся REST-вызовы в системе - там где нужен синхронный ответ прямо сейчас (запрос текущей цены для отображения на странице) - так и остались REST. Не надо тащить Kafka везде где есть HTTP.

Контакт

Нужна такая же инженерная работа?

Опишите задачу и контекст. Ответим в течение рабочего дня, при необходимости подпишем NDA.