Kafka Streams: обработка без кластера

Kafka Streams — не сервер, а библиотека внутри вашего сервиса: state store в RocksDB, changelog-топик как способ пережить падение, exactly-once в обработке. Что это даёт, чем за это платят и когда вместо неё всё-таки нужен Flink

Посчитать агрегат поверх Kafka-топика не значит поднять ещё один кластер. Kafka Streams — это обычная Java-библиотека: new KafkaStreams(topology, config).start() внутри вашего сервиса, без JobManager, TaskManager и отдельного деплоя. Состояние при этом никуда не девается — оно просто живёт не в чужой инфраструктуре, а рядом, в локальном RocksDB, с Kafka в роли журнала, из которого это состояние восстанавливается после падения. Эта статья — про то, как это устроено на живом стенде: топология, state store, changelog-топик и exactly-once, — и чем за модель «библиотека, а не платформа» приходится платить.

Ретрофутуристская схема в духе «Полдня»: одиночный модуль-мастерская подключён к конвейеру-жёлобу Kafka короткой трубой — намеренно один на листе, рядом пустое место с перечёркнутым силуэтом отсутствующего кластера диспетчерских башен. Внутри модуля — картотечный шкаф-хранилище (state store) и катушка-дублёр, повторяющая каждую карточку, уложенную в ящик (changelog). По стене модуля бежит трещина-сбой; катушка отматывается назад и заново заполняет ящик из самой себя, оставаясь целой. Табличка-сравнение противопоставляет одиночный модуль перечёркнутому кластеру башен — «отдельный кластер не нужен»

В статье

Библиотека, а не кластер

Первая статья серии довела заказы из 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 (сумма клиента обнуляется), а из Infinity9.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.20customer_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, для стора с миллионами клиентов — топик, которому требуются те же ресурсы диска и сети, что и любому другому топику кластера, независимо от того, как давно проходила компакция.

Если обработка укладывается в «вход и выход — только Kafka, состояние умеренное, вся логика — деталь одного сервиса», Kafka Streams закрывает задачу без единой лишней сущности в инфраструктуре. Но у модели есть потолок: как только источников и стоков становится много и разнородно (CDC, файлы, БД, lakehouse — не только Kafka), состояние вырастает настолько, что RocksDB и changelog на инстанс сервиса становятся неудобной единицей масштабирования, или нужны savepoint/rescale и управляемые апгрейды параллелизма — это уже задача для отдельной платформы, а не для библиотеки внутри сервиса. Как устроено состояние, чекпоинты и exactly-once во Flink — отдельная статья, а модель времени и окон — в начале серии Flink.

Демо и версии

Стенд — digital-cookbook/messaging/data-streaming: orders.eventsStreamsApp (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 и уровень логирования чувствительны к версии — при обновлении пиновать и перепроверять.

Документация

Обсуждение в Telegram

Присоединиться →

Комментарии