ADG Оставить заявку
Блог Данные и аналитика 4 мин чтения

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 закрывает задачу, не добавляя новых движущихся частей в инфраструктуру.

Контакт

Нужна такая же инженерная работа?

Опишите задачу и контекст. Ответим в течение рабочего дня, при необходимости подпишем NDA.