Обычный 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 в проде — идемпотентность, ребалансировки, мониторинг.
  • Синхронно или асинхронно — когда события вообще нужны.