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

Apache Kafka 2.0 в продакшне: Kafka Streams и EOS наконец делают дубли редкостью

Обновили Kafka-кластер до 2.0 на интеграционном проекте. Exactly Once Semantics в Streams убрал дубли событий при сбоях брокера - то что мешало нам запустить финансовые расчёты на Streams.

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

Apache Kafka 2.0 GA - Kafka Streams improvements и поддержка Exactly Once Semantics (EOS)

Apache Kafka 2.0 вышел в июле, и мы несколько недель смотрели на него со стороны - читали баг-листы, ждали первых патч-версий. В конце августа решили что пора: обновили кластер на одном из интеграционных проектов с 1.1 до 2.0. Главный мотив - Exactly Once Semantics в Kafka Streams. Расскажем почему это было нужно и что получилось.

Контекст: почему Streams, почему сейчас

Проект - финансовая интеграция: несколько источников событий (транзакции, лимиты, корректировки), которые нужно агрегировать в реальном времени и отдавать в downstream-системы. Архитектура на Kafka, consumer-группы на Java, часть агрегации - на самописных батч-скриптах с расписанием.

Kafka Streams как инструмент агрегации мы смотрели ещё на версии 1.0. Подходит органично: данные уже в Kafka, API удобный, KTable для stateful-агрегаций - именно то что нужно. Но был один неприятный сюрприз: при сбое брокера или перебалансировке потребителей возникали дублирующиеся события в выходных топиках. Для мониторинга это терпимо. Для финансовых расчётов, где каждое событие - это деньги, дубль означает неправильный баланс. Поэтому Streams стоял в планах, но с пометкой «когда EOS будет нормально работать».

EOS в Kafka появился ещё в 0.11 - через транзакционный API и idempotent producer. Но полная поддержка в Kafka Streams, со всеми edge-cases при ребалансировке, была сырой. В 2.0 это переработали серьёзно.

Что изменилось в EOS

Ключевое изменение в 2.0 для Streams - новый режим processing.guarantee=exactly_once. До 2.0 он уже существовал, но внутренняя реализация давала нежелательные паузы: при инициализации транзакций каждый StreamThread открывал отдельный Producer. При ребалансировке это создавало ситуации, когда транзакция начата, задача переехала на другой тред, новый тред не знает об открытой транзакции - и брокер ждёт таймаута.

В Kafka 2.0 Streams переключился на так называемый per-task producer: каждая задача имеет своего Producer с уникальным transactional.id, привязанным к задаче, а не к треду. При ребалансировке новый тред, забирающий задачу, инициализирует транзакцию с тем же transactional.id и автоматически прерывает незакрытые транзакции предыдущего. Брокер знает что делать, паузы сокращаются, дублей при сбоях не возникает.

На практике это выглядит так:

Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "financial-aggregator");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka01:9092,kafka02:9092");
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE);
props.put(StreamsConfig.REPLICATION_FACTOR_CONFIG, 3);

Включение буквально в одну строку - EXACTLY_ONCE вместо AT_LEAST_ONCE. Но за этим одним параметром - существенные изменения в поведении брокера и клиента.

Как обновляли кластер

Кластер: три брокера, версия 1.1.1. Обновляли по одному - стандартный rolling upgrade. Kafka поддерживает смешанный кластер на время обновления через inter.broker.protocol.version и log.message.format.version. Выставили эти параметры на 1.1 перед началом, прокатили версию на всех брокерах, потом подняли версии протокола до 2.0.

Несколько моментов которые стоит учесть:

inter.broker.protocol.version=2.0. После обновления всех брокеров нужно не забыть поднять версию протокола. Если оставить на 1.1, часть новых возможностей 2.0 не заработает.

Изменения в consumer group protocol. В 2.0 изменений в протоколе ребалансировки нет - стандартный eager rebalancing работает как раньше.

Deprecation по клиентским API. Убрали несколько устаревших методов Java-клиента. Прошлись по кодовой базе, нашли два места с sendOffsetsToTransaction старой сигнатуры - поправили. Остальное компилировалось без изменений.

Даунтайма не было. Rolling upgrade с тремя брокерами и replication factor=3 прошёл прозрачно для producers и consumers.

Что проверяли после

После обновления и включения EOS в Streams-приложении сделали намеренный сбой: kill -9 на один из брокеров во время активной обработки. До 2.0 в таком сценарии мы наблюдали дубли в выходном топике - consumer downstream видел одно и то же событие дважды. После 2.0 с EXACTLY_ONCE - дублей нет. Kafka восстановила незавершённые транзакции, downstream получил каждое событие ровно один раз.

Дополнительная проверка: искусственная ребалансировка через добавление нового инстанса Streams-приложения. Транзакции корректно перешли на новые таски, паузы в обработке - секунды, не минуты.

Производительность: с EXACTLY_ONCE overhead на throughput есть - транзакционный API добавляет write на __transaction_state топик и дополнительный round-trip для commit. На нашем профиле нагрузки (финансовые события, не высокочастотные) это несущественно. Для потоков с очень высоким throughput стоит мерить отдельно.

Что ещё изменилось в 2.0

Кроме EOS в 2.0 есть несколько полезных изменений:

KStreams API. В Streams-слое улучшена обработка оконных агрегаций и стабилизированы внутренние API. KTable пока эмитирует промежуточные результаты на каждое обновление - буферизация финальных результатов по окнам запланирована в следующих версиях.

AdminClient стал стабильным. До 2.0 AdminClient был @Unstable. Теперь API зафиксирован, и можно управлять топиками программно без опасений что сигнатуры поменяются в следующей версии.

Поддержка Java 11. Kafka 2.0 официально поддерживает Java 11. Java 7 выпилена.

Итог

EOS в Kafka Streams 2.0 - это не маркетинговая фича, а конкретный инженерный результат который разблокировал использование Streams там где at-least-once был неприемлем. Финансовые расчёты на Streams теперь выглядят реалистично, а не как «когда-нибудь потом».

Кластер работает вторую неделю, статистика по дублям - ноль. Будем смотреть дальше: KTable.suppress() открывает новые возможности по агрегациям и мы планируем переписать несколько батч-скриптов на Streams. Если что-то пойдёт не так - расскажем.

Контакт

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

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