Посчитать агрегат поверх Kafka-топика не значит поднять ещё один кластер. Kafka Streams — это обычная Java-библиотека: new KafkaStreams(topology, config).start() внутри вашего сервиса, без JobManager, TaskManager и отдельного деплоя. Состояние при этом никуда не девается — оно просто живёт не в чужой инфраструктуре, а рядом, в локальном RocksDB, с Kafka в роли журнала, из которого это состояние восстанавливается после падения. Эта статья — про то, как это устроено на живом стенде: топология, state store, changelog-топик и exactly-once, — и чем за модель «библиотека, а не платформа» приходится платить.
В статье
- Библиотека, а не кластер
- Топология: от заказа до агрегата
- State store и RocksDB
- Changelog как механизм восстановления
- Exactly-once и его граница
- Цена модели
- Когда всё-таки Flink
- Демо и версии
- Документация
Библиотека, а не кластер
Первая статья серии довела заказы из PostgreSQL до топика `orders.events` — обычный поток событий в Kafka, без своей семантики агрегации. Дальше нужно посчитать, сколько заказов и на какую сумму сделал каждый клиент, и держать этот агрегат актуальным. Можно поднять для этого Flink-кластер — а можно добавить в сервис, который и так уже говорит с Kafka, зависимость на kafka-streams и написать топологию прямо в коде.
Разница не в мощности, а в операционной модели. Flink-задание — это dataflow-граф, который планируется и исполняется отдельным рантаймом (JobManager/TaskManager); что это такое и как там устроены время и состояние — отдельная серия, здесь не пересказывается. Kafka Streams — топология операторов, которая исполняется потоками (StreamThread) внутри процесса вашего сервиса: сколько инстансов сервиса — столько параллелизма, никакого отдельного планировщика ресурсов не появляется. Это и удобство («ничего нового не разворачивать»), и ограничение («вход и выход — только Kafka, масштабирование — только через инстансы сервиса») одновременно.
Топология: от заказа до агрегата
Стенд (digital-cookbook/messaging/data-streaming) считает по заказу из orders.events бегущий агрегат по клиенту и пишет его в customer.totals. Топология собрана сырым Processor API (не DSL) — это осознанный выбор темы статьи, и вот почему это важно, прямо из комментариев StreamsApp.java:
public class StreamsApp {
static final String IN = "orders.events";
static final String OUT = "customer.totals";
static final String STORE = "customer-totals-store";
public static void main(String[] args) {
Properties p = new Properties();
p.put(StreamsConfig.APPLICATION_ID_CONFIG, "ds-streams");
p.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, env("KAFKA", "localhost:9096"));
p.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
p.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
p.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);
p.put(StreamsConfig.STATE_DIR_CONFIG, env("STATE_DIR", "/tmp/ds-streams"));
StreamsBuilder b = new StreamsBuilder();
b.addStateStore(Stores.keyValueStoreBuilder(
Stores.persistentKeyValueStore(STORE), Serdes.String(), Serdes.String()));
b.<String, String>stream(IN, Consumed.with(Serdes.String(), Serdes.String()))
// Ключ orders.events — order_id (Debezium: aggregateid outbox-события).
// Processor API САМ ПО СЕБЕ не меняет партиционирование: addStateStore()
// просто прикрепляет стор к тем же task'ам, что читают входные партиции.
// Пока orders.events была 1 партиция, все заказы физически проходили через
// единственный task — баг маскировался. При нескольких партициях заказы
// одного клиента с разными order_id разъедутся по разным task'ам с
// ИЗОЛИРОВАННЫМИ сторами, и агрегация по клиенту станет молча неверной
// (несколько частичных сумм на одного клиента вместо одной полной).
// Поэтому переключаем ключ на customer_id и ЯВНО репартиционируем: Kafka
// Streams создаст промежуточный repartition-топик и перетасует записи так,
// чтобы все события одного клиента шли в один и тот же task/партицию/стор
// (co-partitioning). DSL-агрегация (groupBy) сделала бы это неявно сама —
// здесь это делаем руками, потому что тема статьи — сырой Processor API
// со state store.
.flatMap((orderId, json) -> {
ParsedEvent ev = parseEvent(json);
if (ev == null) {
System.err.println("пропущено битое событие (order=" + orderId + "): невалидный JSON, нечисловой/отсутствующий customer_id или amount, либо amount не конечен (NaN/Infinity)");
return List.<KeyValue<String, String>>of();
}
return List.of(KeyValue.pair(Long.toString(ev.customerId()), Double.toString(ev.amount())));
})
// Число партиций repartition-топика ниже НЕ задано явно (нет
// .withNumberOfPartitions(N)) — Kafka Streams наследует его от апстрима,
// то есть от фактического числа партиций orders.events на момент первого
// запуска приложения. Это неявная зависимость: пересоздание orders.events
// с другим числом партиций молча меняет и число партиций
// ds-streams-customer-totals-repartition вместе с ним.
.repartition(Repartitioned.<String, String>as("customer-totals")
.withKeySerde(Serdes.String())
.withValueSerde(Serdes.String()))
.process(TotalsProcessor::new, STORE)
.to(OUT, Produced.with(Serdes.String(), Serdes.String()));
KafkaStreams streams = new KafkaStreams(b.build(), p);
// Политика одна и БЕЗУСЛОВНАЯ: любая необработанная ошибка stream-треда
// останавливает клиента. REPLACE_THREAD не используется нигде, поэтому
// и разбора типов исключений здесь нет — ответ не зависит от причины.
// Обоснование выбора — в разделе про changelog ниже.
streams.setUncaughtExceptionHandler(throwable -> {
System.err.println("необработанная ошибка stream-треда: " + throwable);
return StreamsUncaughtExceptionHandler.StreamThreadExceptionResponse.SHUTDOWN_CLIENT;
});
streams.start();
}
}Поток: .stream(IN) читает orders.events → .flatMap(...) разбирает и валидирует JSON, оставляет только пару (customer_id, amount), отбрасывая всё битое → .repartition(...) перетасовывает записи так, чтобы события одного клиента всегда оказывались в одном таске → .process(TotalsProcessor::new, STORE) читает и обновляет агрегат в state store → .to(OUT) пишет результат в customer.totals.
Одна деталь валидации стоит отдельного упоминания, потому что на ней легко проглядеть тихую порчу данных. Double.parseDouble принимает "NaN", "Infinity" и "-Infinity" — это валидные литералы Java, а не ошибка разбора: NumberFormatException на них не бросается, и общий catch их не поймает. Дальше такое значение доехало бы до агрегата и исказило бы его молча — на округлении Math.round(total * 100.0) / 100.0 из NaN получается ровно 0.0 (сумма клиента обнуляется), а из Infinity — 9.223372036854776E16, потому что Math.round от бесконечности даёт Long.MAX_VALUE. Это не отказ, который кто-то заметит по логам, а неверное число в витрине. Поэтому после разбора стоит явная проверка Double.isFinite(amount), и такое событие отбрасывается наравне с невалидным JSON. Симметричная проверка Double.isFinite(total) стоит и после суммирования — накопленная сумма может переполнить double даже из конечных слагаемых, и тогда обработка останавливается громко, а не пишет то же 9.2e16.
Целые числа прячут ту же ловушку, только последствия хуже. JsonNode.isIntegralNumber() истинен и для BigInteger, который не влезает в long, а asLong() в этом случае молча усекает: customer_id=18446744073709551617 (это 2^64+1) превращается в asLong()==1. То есть событие с битым идентификатором приписалось бы реальному клиенту 1 — чужая сумма тихо смешалась бы с настоящей. Поэтому перед asLong() нужен canConvertToLong(). Тот же принцип на выходе: инкремент счётчика заказов сделан через Math.addExact, потому что обычное +1 на Long.MAX_VALUE даёт Long.MIN_VALUE и записало бы в changelog отрицательное число заказов. Общее правило простое: у любого «числа из внешних данных» есть значения, которые проходят разбор, но ломают смысл, — и почти всегда они превращаются не в ошибку, а в правдоподобное неверное число.
«Отбрасывая» здесь — по-настоящему теряя, а не откладывая. flatMap печатает причину в stderr (строка — прямо в коде выше) и возвращает пустой список: исходящей записи не появляется нигде, ни в customer.totals, ни в отдельном топике. Никакого DLQ на этом шаге нет — DLQ в этой серии живёт на стоке, у Go-консьюмера на выходе из customer.totals (ст.3 серии), а не у Kafka Streams на входе. Границу стоит держать в голове честно: битое входное событие Kafka Streams отбрасывает молча — со строкой в stderr, но без записи, которую потом можно было бы разобрать.
Ключевой тезис, который стоит унести в прод: явный .repartition() здесь не украшение, а необходимость. DSL-метод groupBy() вставил бы репартиционирование сам, незаметно. Сырой Processor API (addStateStore() + .process()) этого не делает — он просто прикрепляет стор к тем таскам, что и так читают входные партиции. Если пропустить .repartition() при входной топологии больше чем с одной партицией, агрегация по клиенту станет молча неверной: заказы одного клиента с разными order_id разъедутся по разным таскам с изолированными сторами, и вместо одного полного агрегата получится несколько частичных. Это ровно тот шов между DSL и Processor API, который в комментариях к коду и назван причиной, почему репартиционирование сделано руками, а не понадеялось на библиотеку.
State store и RocksDB
Stores.persistentKeyValueStore(STORE) — это персистентный key-value store поверх embedded RocksDB, один инстанс на таск, физически на диске по пути STATE_DIR_CONFIG (/tmp/ds-streams в стенде). TotalsProcessor в process() делает с ним ровно то же, что сделал бы с любой локальной map: store.get(cust), посчитал новый orders/total, store.put(cust, json). Никакого сетевого похода за состоянием — стор живёт в том же процессе, что и логика агрегации, поэтому чтение и запись быстрые.
Плата за эту локальность — стор привязан к конкретному инстансу и конкретному диску. Если инстанс с этим таском умирает и не поднимается заново на том же диске, локальные файлы RocksDB для этого таска пропадают вместе с ним. Обычная key-value БД на этом бы и закончилась — нужен был бы отдельный бэкап. У Kafka Streams эта проблема закрыта на уровне библиотеки: каждая запись в стор одновременно уходит в changelog-топик, обычный топик Kafka, который переживает падение инстанса, потому что живёт в Kafka, а не на локальном диске.
Changelog как механизм восстановления
Для STORE=customer-totals-store библиотека автоматически заводит топик ds-streams-customer-totals-store-changelog, наследующий число партиций от входа через co-partitioning:
$ docker exec ds-kafka /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --list
...
ds-streams-customer-totals-repartition
ds-streams-customer-totals-store-changelog
...
$ docker exec ds-kafka /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic ds-streams-customer-totals-store-changelog
Topic: ds-streams-customer-totals-store-changelog PartitionCount: 3 ReplicationFactor: 1 Configs: min.insync.replicas=1,cleanup.policy=compact,message.timestamp.type=CreateTime
cleanup.policy=compact — не случайность: changelog хранит не историю событий, а «последнее значение по ключу», ровно как содержимое стора. После обработки первых 50 заказов (5 клиентов × 10 заказов каждый) в топике — по одной записи на каждое обновление состояния клиента, всего 50:
2 {"customer_id":2,"orders":8,"total":384.32}
3 {"customer_id":3,"orders":7,"total":388.52}
2 {"customer_id":2,"orders":9,"total":474.39}
3 {"customer_id":3,"orders":8,"total":413.0}
3 {"customer_id":3,"orders":9,"total":472.04}
3 {"customer_id":3,"orders":10,"total":521.99}
2 {"customer_id":2,"orders":10,"total":565.71}
4 {"customer_id":4,"orders":1,"total":8.83}
...
4 {"customer_id":4,"orders":9,"total":466.99}
4 {"customer_id":4,"orders":10,"total":492.2}
Processed a total of 50 messages
PG в этот момент подтверждает то же самое: SELECT customer_id, count(*), sum(amount) FROM orders GROUP BY customer_id даёт 4|10|492.20 — customer_id=4 в changelog (orders=10, total=492.2) совпадает с PG до цента.
Проверка восстановления на стенде — жёсткая: процесс убит kill -9 (не graceful shutdown), локальный каталог /tmp/ds-streams удалён полностью (rm -rf), приложение перезапущено той же командой (KAFKA=localhost:9096 mvn -q compile exec:java).
Честно: строка вида «restoring state store» в stdout/stderr на INFO-уровне (дефолтный порог slf4j-simple) не появилась — при таком объёме (50 записей) восстановление в этой версии Kafka Streams укладывается ниже порога логирования. Вместо строки лога — доказательство сильнее и фальсифицируемое, физическое и функциональное:
$ ls /tmp/ds-streams # ДО перезапуска: rm -rf стёр каталог
ls: cannot access '/tmp/ds-streams': No such file or directory
$ cd streams && KAFKA=localhost:9096 mvn -q compile exec:java # рестарт
$ sleep 8; du -sh /tmp/ds-streams # сразу после рестарта
4.1M /tmp/ds-streams
$ sleep 10; du -sh /tmp/ds-streams # чуть позже
13M /tmp/ds-streams
$ find /tmp/ds-streams -path "*rocksdb*" | head -5
/tmp/ds-streams/ds-streams/1_0/rocksdb
/tmp/ds-streams/ds-streams/1_0/rocksdb/customer-totals-store
/tmp/ds-streams/ds-streams/1_0/rocksdb/customer-totals-store/LOG
/tmp/ds-streams/ds-streams/1_0/rocksdb/customer-totals-store/000004.log
RocksDB-каталог заново материализовался с данными — на пустом состоянии его бы не было: на 8-й секунде после рестарта — уже 4.1M (восстановление ещё в процессе), на 18-й — 13M, тот же итоговый размер, что и до kill -9. Функциональная проверка ещё жёстче: до рестарта customer_id=1 был на orders=10, total=636.70 (та же цифра, что и в PG-проверке чуть выше). Один новый заказ на 10.00 после рестарта:
$ docker exec ds-kafka /opt/kafka/bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic customer.totals --from-beginning --timeout-ms 8000 | grep '"customer_id":1,'
{"customer_id":1,"orders":1,"total":80.33}
...
{"customer_id":1,"orders":10,"total":636.7}
{"customer_id":1,"orders":11,"total":646.7}
orders=11, total=646.70 (636.70 + 10.00) — счётчик продолжил накопление с 10, а не сбросился на 1. Это и есть падающий вариант, которым проверка фальсифицируема: если бы состояние не восстановилось из changelog (осталось бы пустым), новая запись показала бы orders=1, total=10.00. Отсутствие строки лога здесь не дефект стенда — Kafka Streams, как и Debezium в статье про CDC, на INFO не логирует подобные вещи при небольшом объёме; доказательство через содержимое RocksDB-каталога и продолжившийся счётчик надёжнее одной строки лога именно потому, что оно фальсифицируемо, а строка лога — нет.
Важно не путать это с другим сценарием — повреждённой ЗАПИСЬЮ внутри самого стора, а не отсутствием стора целиком, как выше. TotalsProcessor.process() на каждое обновление сначала читает store.get(cust): если значение не null, но не разбирается как JSON или не содержит полей orders/total, код не пытается молча продолжить — он бросает выделенное исключение CorruptedStateStoreException, и приложение останавливается: чистая однократная остановка, которую оператор увидит.
Это осознанный выбор между тремя вариантами. Тихий сброс в orders=1 (как было в ранней редакции стенда) — худший: накопленная история клиента теряется навсегда, заниженное значение уходит в changelog и дальше в customer.totals/ClickHouse без коррекции, и ни один следующий шаг конвейера этого не заметит — данные молча испорчены. REPLACE_THREAD здесь неверен: причина детерминированная (та же повреждённая запись), поток поднимался бы и падал на ней бесконечно — краш-луп без движения. SHUTDOWN_CLIENT останавливает приложение один раз и громко — это не потеря данных, а сигнал: локальный стор в неизвестном состоянии, дальше идти нельзя.
Отсюда и вид самого обработчика: он безусловный, без разбора типа исключения. Соблазн написать «для повреждения стора — SHUTDOWN_CLIENT, для остального — REPLACE_THREAD» выглядит аккуратнее, но перезапуск потока безопасен только для причин, которые вы явно классифицировали как транзиентные. Неизвестное исключение с тем же успехом окажется детерминированным багом — NullPointerException на конкретном входе, — и REPLACE_THREAD даст ровно тот же краш-луп. Код этого стенда транзиентных причин не порождает, поэтому классификации нет вовсе: писать ветвление, у которого все ветки возвращают один ответ, значило бы имитировать выбор, которого не существует. Разбор типов остаётся в тестах — там он проверяет, что упало именно то исключение с внятным сообщением (и там же приходится идти по цепочке getCause(): Kafka Streams оборачивает исключение из process() в StreamsException, так что верхний объект — не то, что бросил ваш код).
Честный учебный вывод: правильный ответ на повреждённую запись стора в проде — именно fail-fast, после которого оператор пересоздаёт state store и делает replay состояния из того же changelog-топика (истина там уже есть — её восстановление и показано чуть выше, просто локальный RocksDB её потерял или испортил), либо явная миграция формата значения, если повреждение вызвано изменившейся схемой, а не физическим сбоем. Companion-код учебной статьи не должен нести в себе тихо портящий данные fallback — управляемый простой на восстановление дешевле молчаливо заниженного агрегата. На самом стенде этот путь недостижим (стор пишет и читает один и тот же код, формат не расходится), поэтому fail-fast не меняет ни одного числа демонстраций — он просто делает пример корректным по умолчанию; проверяется он отдельным TopologyTestDriver-тестом (см. streams/src/test).
Exactly-once и его граница
PROCESSING_GUARANTEE_CONFIG=EXACTLY_ONCE_V2 — транзакционный режим: чтение из orders.events, обновление state store и запись в customer.totals коммитятся атомарно, в одной транзакции Kafka. Внутреннее устройство этого механизма — транзакционные продюсеры, идемпотентность, координатор транзакций — разбирает статья про exactly-once и транзакции в Kafka; здесь важна граница, а не механика.
Граница простая и явно зафиксирована в комментарии к классу StreamsApp: EXACTLY_ONCE_V2 покрывает read-process-write внутри Kafka — от чтения orders.events до записи в customer.totals. Дальше по цепочке стоит Go-консьюмер, который читает customer.totals и пишет в ClickHouse — это уже отдельный шаг, вне транзакции Kafka Streams. Там дедуп и идемпотентность — своя забота стока, а не то, что «уже гарантировано» exactly-once Kafka Streams. Тот же принцип — гарантия держится на границах, а не переносится по цепочке сама собой — во Flink устроен иначе (чекпоинты и two-phase commit к стокам), но исходная граница «внутри платформы vs за её пределами» одна и та же для обоих движков.
Цена модели
Модель «библиотека вместо кластера» не бесплатна — платят за неё по трём осям.
Восстановление состояния — не мгновенное. RocksDB-каталог на 50 записей агрегата занял 13M и в стенде уже достиг этого размера к 18-й секунде после kill -9 и полного удаления state dir; раньше, на 8-й секунде, каталог был на 4.1M — восстановление ещё шло. Это верхняя оценка по двум контрольным точкам, а не точно измеренная длительность. Цифра сама по себе невелика, но масштабируется с размером стора: чем больше уникальных ключей и чем реже они компактятся, тем дольше таск читает свой changelog с начала при холодном старте на новом инстансе.
Ребаланс двигает таски вместе с их состоянием. Если таск с состоянием переезжает на другой инстанс (масштабирование, деплой, падение соседа), новому инстансу нужно либо иметь под рукой standby-копию, либо восстановить стор с нуля из changelog — по тому же механизму, что показан выше. Восстановление здесь не разовая демонстрация, а штатная процедура, которая срабатывает при каждом таком переезде.
Changelog — это ещё один топик Kafka, который нужно эксплуатировать. cleanup.policy=compact в теории тянет размер топика к нижней границе «одна актуальная запись на уникальный ключ» (а не растёт с числом событий), но компакция асинхронна: она проходит по сегментам с задержкой, не после каждой записи, и старые версии значения по тому же ключу могут какое-то время оставаться на диске рядом с новыми. Строгой пропорциональности размера числу уникальных ключей в моменте нет — до очередного прохода компакции changelog может быть заметно больше своего теоретического минимума, и это всё равно данные, реплицируемые и хранимые Kafka наравне с прикладными топиками. Для стора с 50 записями это ds-streams-customer-totals-store-changelog с PartitionCount: 3, для стора с миллионами клиентов — топик, которому требуются те же ресурсы диска и сети, что и любому другому топику кластера, независимо от того, как давно проходила компакция.
Когда всё-таки Flink
Если обработка укладывается в «вход и выход — только Kafka, состояние умеренное, вся логика — деталь одного сервиса», Kafka Streams закрывает задачу без единой лишней сущности в инфраструктуре. Но у модели есть потолок: как только источников и стоков становится много и разнородно (CDC, файлы, БД, lakehouse — не только Kafka), состояние вырастает настолько, что RocksDB и changelog на инстанс сервиса становятся неудобной единицей масштабирования, или нужны savepoint/rescale и управляемые апгрейды параллелизма — это уже задача для отдельной платформы, а не для библиотеки внутри сервиса. Как устроено состояние, чекпоинты и exactly-once во Flink — отдельная статья, а модель времени и окон — в начале серии Flink.
Демо и версии
Стенд — digital-cookbook/messaging/data-streaming: orders.events → StreamsApp (Processor API, EXACTLY_ONCE_V2) → customer.totals, с явным репартиционированием и state store в RocksDB. docker-compose, scripts/seed.sh для генерации заказов и шаги воспроизведения восстановления state store — в README демо.
Версии: Kafka 4.3.1 (образ apache/kafka:4.3.1), Kafka Streams 4.3.1 — поставляется вместе с Kafka, версия совпадает с версией брокера (streams/pom.xml: <kafka.version>4.3.1</kafka.version>); сборка — Jackson (databind) 2.18.2, JDK Temurin 21.0.11, Maven 3.9.16. Поведение восстановления state store и уровень логирования чувствительны к версии — при обновлении пиновать и перепроверять.
Документация
- Официально: Kafka Streams — docs, Processor API, exactly-once (KIP-129/447).
- Серия Flink: модель, event-time, watermarks и окна, состояние и exactly-once во Flink.
- Смежное на сайте: сборка пайплайна PostgreSQL → Debezium → Kafka → ClickHouse (ст.1 серии), проектирование событийных пайплайнов (ст.3 серии), exactly-once и транзакции в Kafka, хаб «Messaging: выбрать и эксплуатировать».
Комментарии