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.