Kafka Streams как альтернатива Spark Streaming: потоковая обработка без отдельного кластера
Kafka Streams - встроенная библиотека потоковой обработки, анонсированная в 2016. Сравниваем с batch-ETL и Spark Streaming: когда хватает библиотеки без кластера.
Kafka Streams API анонсирован в 2016 году как встроенная библиотека потоковой обработки без внешнего кластера обработки
На проекте, где мы пробовали Kafka как шину событий, пилот постепенно перерос в рабочий контур. Брокер стоит, данные текут, consumer пишет в DWH. И тут заказчик спросил: а можно прямо в потоке считать скользящие суммы по клиентам, не нагружая хранилище? Вопрос разумный: тащить сырые события в DWH, агрегировать там - это один подход. Делать агрегацию до записи - другой.
Смотреть на Spark Streaming мы начали первым делом - он у всех на слуху. Но потом в апрельских анонсах от Confluent увидели Kafka Streams, который идёт в составе самой Kafka как обычная Java-библиотека. Без отдельного кластера обработки. Без YARN, без Mesos, без standalone-мастера. Это стоило пощупать.
Что такое Kafka Streams на практике
Kafka Streams - это не система типа Spark или Flink, которую надо разворачивать как отдельную службу. Это библиотека: добавляешь зависимость в Maven, пишешь обычное Java-приложение, оно само читает из топиков, обрабатывает, пишет в выходные топики. Параллелизм и масштабирование достигаются запуском нескольких экземпляров одного приложения - каждый берёт часть партиций.
Архитектурно это выглядит примерно так:
input topic -> [Kafka Streams App] -> output topic -> consumer -> DWH
В самом приложении описываешь KStream или KTable - потоки и таблицы-состояния - и операции над ними: фильтрация, маппинг, агрегация с окнами. Состояние (например, накопленные суммы по окну) хранится локально через RocksDB - встраиваемая key-value база, никакого внешнего хранилища состояния не нужно.
Пример того, что мы пробовали - подсчёт оборота по клиентам за последние 5 минут скользящим окном:
KStream<String, Transaction> txStream = builder.stream("transactions");
KTable<Windowed<String>, Long> clientTotals = txStream
.groupBy((key, tx) -> tx.getClientId())
.windowedBy(TimeWindows.of(5 * 60 * 1000L))
.aggregate(
() -> 0L,
(clientId, tx, total) -> total + tx.getAmount(),
Serdes.Long(),
"client-totals-store"
);
clientTotals.toStream().to("client-totals-output");
Результат - в отдельный топик, откуда его читает лёгкий consumer и пишет в таблицу DWH. Аналитики видят скользящие агрегаты с задержкой в единицы секунд.
Почему не Spark Streaming
Spark Streaming мы тоже смотрели, и он умеет больше. Но цена этого «больше» - отдельный Spark-кластер. Для нашего случая это значит: дополнительные ноды, Hadoop или standalone-режим, отдельный операционный контур с мониторингом JVM на нескольких машинах, отдельная команда компетенций.
На проектах, где уже есть Hadoop-кластер и нужна сложная потоковая аналитика - Spark Streaming выглядит органично. У нас Hadoop нет, задача - несколько агрегаций над одним потоком. Разворачивать ради этого кластер кажется избыточным.
Kafka Streams вписывается в то, что уже есть: Kafka-брокеры стоят, Java-разработчик в команде есть. Ещё одно приложение рядом с consumer - операционно это ничем не отличается от того, что уже есть.
Сравнение с batch-ETL
Наш исходный ночной ETL работал по схеме: выгрузка за сутки, трансформация, загрузка. Аналитика по вчерашним данным. Это нормально для отчётов, которые смотрят утром раз в день. Для оперативного мониторинга - нет.
Kafka Streams в связке с брокером меняет это: агрегации обновляются непрерывно. Но надо понимать ограничения:
- Состояние - локально. RocksDB живёт на той машине, где запущено приложение. При падении Kafka Streams восстанавливает состояние из changelog-топика, но это занимает время. Для больших окон (часы, дни) восстановление может быть медленным.
- Ровно-один-раз - сложно. At-least-once из коробки, exactly-once - требует поддержки транзакций на уровне брокера, которой в текущей версии Kafka нет. Для финансовых агрегатов это значит: думать про идемпотентность на стороне DWH.
- Только Kafka. Источник и приёмник - только топики Kafka. Это не проблема, если Kafka уже центральная шина. Если нет - придётся сначала выстраивать её.
- Мониторинг. Инструментирование приложения через JMX, метрики задержки и лага доступны, но мониторинг надо строить самим - готового дашборда нет.
Где мы сейчас
Приложение на Kafka Streams запущено рядом с основным consumer на тех же VM. Деплой - обычный jar через systemd, никакой экзотики. Скользящие агрегаты пишутся в отдельные таблицы DWH, аналитики смотрят их через существующие отчёты.
На проекте по DWH и аналитике это первый опыт потоковой обработки без отдельной инфраструктуры. Пока выглядит как разумный компромисс: не платишь ценой сложности за то, что нужно относительно просто.
Если задача вырастет до join-ов нескольких потоков или сложных паттернов обнаружения - нужно будет смотреть шире. Но для агрегаций одного потока с окнами Kafka Streams закрывает задачу, не добавляя новых движущихся частей в инфраструктуру.