Обычный consumer читает сообщения и делает что-то с каждым по отдельности. Но часть задач в эту модель не помещается: посчитать заказы за последние пять минут, склеить поток платежей с потоком заказов, держать актуальный остаток по каждому счёту. Это потоковая обработка — и Kafka Streams делает её обычной Java-библиотекой: без отдельного кластера, внутри вашего сервиса.
Что это такое
Kafka Streams — библиотека из состава Apache Kafka: добавляете зависимость, описываете топологию (из каких топиков читать, как преобразовывать, куда писать) — и приложение запускается как обычный Spring Boot сервис. Сравните с Flink или Spark Streaming, которым нужен собственный кластер с менеджментом задач. Масштабирование — то же, что у consumer groups: запустили второй экземпляр сервиса — партиции поделились между ними.
StreamsBuilder builder = new StreamsBuilder();
builder.stream("orders", Consumed.with(Serdes.String(), orderSerde))
.filter((key, order) -> order.amount().compareTo(BigDecimal.ZERO) > 0)
.groupByKey()
.count()
.toStream()
.to("order-counts");
KStream и KTable: событие против состояния
Главная пара понятий:
- KStream — поток событий: каждое сообщение — самостоятельный факт («заказ создан», «платёж прошёл»). Новое сообщение ничего не отменяет.
- KTable — таблица-состояние: сообщения с одним ключом перезаписывают друг друга, таблица хранит последнее значение («текущий остаток счёта», «актуальный профиль»).
Одни и те же данные можно читать обоими способами — вопрос семантики: история изменений или текущее состояние. Compacted-топик — естественная пара к KTable.
Stateful-операции: окна, join, state stores
Сила Kafka Streams — операции с состоянием:
- Агрегации по окнам: «количество заказов за 5 минут» — окно едет по времени, счётчики живут в state store (встроенный RocksDB на диске экземпляра).
- Join потоков: склеить платёж с заказом по ключу в пределах окна — классическая задача «ждём второе событие пары».
- Отказоустойчивость состояния: каждый state store дублируется в changelog-топик Kafka; упал экземпляр — новый восстановит состояние из топика и продолжит.
Это то, что руками поверх consumer писать долго и ошибочно: своё хранилище, своя логика восстановления, свои дедлайны окон.
Exactly-once
Kafka Streams поддерживает режим processing.guarantee=exactly_once_v2: чтение, обработка, запись результата и сдвиг offset объединяются в транзакцию Kafka. Для конвейеров «топик → обработка → топик» это честный exactly-once без самодельной идемпотентности. Важная граница: гарантия действует внутри мира Kafka — вызов внешнего REST API или запись в базу транзакцией не покрываются, там идемпотентность остаётся вашей заботой.
Когда Kafka Streams, а когда нет
Берите, если: агрегации по времени, join потоков, поддержание производного состояния, материализация «текущего вида» из событий — и данные уже в Kafka.
Не берите, если:
- обработка — «прочитал сообщение, вызвал сервис, записал в базу»: хватит обычного consumer, Streams добавит сложности без пользы;
- нужны тяжёлые вычисления на большом кластере, SQL по потокам, ML-конвейеры — это территория Flink;
- данных нет в Kafka: Streams читает и пишет только Kafka-топики.
Коротко
- Kafka Streams — библиотека потоковой обработки: топология внутри вашего сервиса, без отдельного кластера; масштабирование — партициями, как у consumer groups.
- KStream — поток независимых событий, KTable — состояние «последнее значение по ключу»; выбор — вопрос семантики данных.
- Stateful-операции (окна, join, агрегации) хранят состояние в локальном RocksDB с восстановлением из changelog-топика.
- exactly_once_v2 даёт честный exactly-once внутри Kafka; внешние вызовы и базы — по-прежнему через идемпотентность.
- Простая обработка «сообщение → действие» — обычный consumer; тяжёлая аналитика — Flink; производное состояние и окна над Kafka-данными — Streams.
Что почитать дальше
- Основы Kafka — партиции, offset, consumer groups: фундамент, на котором стоит Streams.
- Kafka в проде — идемпотентность, ребалансировки, мониторинг.
- Синхронно или асинхронно — когда события вообще нужны.