KSQL на Kafka: аналитики пишут SQL, Java Streams лежат на полке
Подняли Confluent KSQL поверх Kafka с POS-событиями: real-time агрегации без Java, аналитики получили привычный SQL и первые инсайты за часы вместо дней.
Confluent KSQL - SQL-интерфейс для потоковой обработки поверх Kafka, активно развивается в 2019 году
В феврале мы разбирали ClickHouse с materialized views для аналитики событий. Там источником были пользовательские клики, и Kafka выступала транспортом. Этот пост про другой проект - торговая сеть, POS-терминалы, и вопрос «что происходит прямо сейчас» вместо «что было вчера».
Откуда задача
Клиент - розничная сеть с несколькими сотнями точек продаж. Каждый POS-терминал при пробитии чека пишет событие в Kafka: товары, количество, сумма, кассир, касса, магазин. Поток нешуточный - несколько тысяч событий в минуту в пиковые часы.
До нашего появления аналитика жила в классическом DWH: события из Kafka забирал Python-скрипт, складывал в PostgreSQL, ночью отдельный ETL-процесс строил агрегаты, утром аналитик смотрел на вчерашние данные. Если продажи в конкретном магазине проседали в обед - об этом узнавали завтра. Если кассир по ошибке дважды пробивал дорогостоящий товар - тоже завтра.
Запрос был простой: хотим видеть продажи по магазинам и категориям с задержкой минуты, а не сутки. Желательно без найма Java-разработчика.
Почему не Kafka Streams
Очевидный ответ на «real-time агрегации из Kafka» в 2019 году - это Kafka Streams. Библиотека, встраивается в приложение, всё красиво. Мы на неё смотрели.
Проблема оказалась не техническая. У клиента нет Java-разработчиков. Есть аналитики, которые умеют SQL, и есть один Python-разработчик, который поддерживает скрипты. Писать и поддерживать Java Streams в таком контексте - это найм нового человека или аутсорс на постоянной основе. Нам нужно было что-то, что аналитики смогут читать и понимать сами.
Confluent KSQL решал именно эту задачу. SQL поверх Kafka Streams - снизу та же библиотека, сверху синтаксис, который аналитик читает без словаря.
Как выглядит KSQL на практике
KSQL работает со стримами и таблицами. Стрим - это бесконечная последовательность событий из топика Kafka. Таблица - это агрегированное состояние на текущий момент, по сути материализованный результат GROUP BY.
Сначала объявили стрим поверх топика с чековыми событиями:
CREATE STREAM pos_events (
receipt_id VARCHAR,
store_id VARCHAR,
cashier_id VARCHAR,
category VARCHAR,
item_id VARCHAR,
quantity INT,
amount DOUBLE,
event_ts BIGINT
) WITH (
KAFKA_TOPIC='pos-receipts',
VALUE_FORMAT='JSON',
TIMESTAMP='event_ts'
);
Дальше - агрегирующая таблица с продажами по магазинам за скользящие 15 минут:
CREATE TABLE store_sales_15m AS
SELECT
store_id,
category,
COUNT(*) AS receipt_count,
SUM(amount) AS total_amount,
SUM(quantity) AS total_items
FROM pos_events
WINDOW TUMBLING (SIZE 15 MINUTES)
GROUP BY store_id, category;
Это всё. Аналитик смотрит на этот SQL и сразу понимает что происходит. Никакого Java, никакого DSL для Streams API. Результат живёт в отдельном Kafka-топике, откуда его можно забирать в Grafana или куда угодно ещё.
Для оперативного мониторинга добавили вторую таблицу - почасовую сводку по кассирам:
CREATE TABLE cashier_hourly AS
SELECT
cashier_id,
store_id,
COUNT(*) AS receipts,
SUM(amount) AS revenue,
AVG(amount) AS avg_receipt
FROM pos_events
WINDOW TUMBLING (SIZE 1 HOUR)
GROUP BY cashier_id, store_id;
Что получилось по факту
Время до первого инсайта. Раньше, чтобы добавить новую аналитическую метрику, аналитик формулировал задачу, разработчик писал ETL-логику, тестировал, деплоил - это дни. Теперь аналитик сам написал новый CREATE TABLE в KSQL CLI, проверил что данные идут, и метрика появилась в потоке. Несколько часов вместо нескольких дней - и это не преувеличение.
Задержка данных. От события на кассе до появления в агрегате - около минуты. Это Kafka latency плюс время до закрытия окна агрегации. Для бизнес-мониторинга продаж - более чем достаточно.
Аномалии. На второй день после запуска аналитики сами, без нас, написали запрос на обнаружение нехарактерно высоких средних чеков по конкретным кассирам и нашли один кейс с похожим на ошибку двойным пробитием. Раньше такое искали в сверках через неделю.
Где KSQL не идеален
Windowing семантика. Tumbling-окна замечательно работают для агрегации «за последние N минут», но если нужен скользящий результат - например, «продажи за любые последние 15 минут с перекрытием» - это Hopping Windows с другой логикой перекрытия. Аналитикам пришлось объяснять разницу, и не один раз.
Late events. Если событие с терминала пришло с задержкой (нестабильный интернет в магазине - реальная ситуация), оно может попасть в закрытое окно. KSQL поддерживает grace period для опоздавших событий, но это надо настраивать осознанно. По умолчанию опоздавшее событие просто теряется для уже закрытого окна.
Мониторинг самого KSQL. Метрики выставляются через JMX, интеграция с Prometheus - через отдельный JMX exporter. Не страшно, но ещё одна движущаяся часть в инфраструктуре.
Persistent queries живут в KSQL Server. Если сервер упал - все запущенные агрегирующие таблицы останавливаются. При рестарте они поднимаются и начинают читать Kafka с последнего оффсета, исторические данные за время простоя в новые окна не попадают. Это надо учитывать при оценке надёжности.
Промежуточный итог
KSQL закрыл конкретную задачу: real-time агрегации из Kafka без Java-разработчика. Аналитики пишут запросы сами, итерируют быстро, не ждут разработки ETL. Python-разработчик клиента занимается подключением результирующих топиков к дашбордам, а не написанием потоковой логики.
Это не замена Kafka Streams для сложных кейсов - там где нужны кастомные десериализаторы, сложный stateful processing или вызовы внешних сервисов, Java Streams гибче. Но для аналитических агрегаций в командах без Java-разработчиков KSQL работает. Подробнее о таких проектах - в разделе DWH и аналитика.