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

Apache Kafka 0.9 как шина событий: пробуем заменить ночной ETL на pub/sub

Kafka 0.9 вышла с новым consumer API и Kerberos-аутентификацией. Пробуем её как шину между ERP и аналитическим хранилищем вместо ночного файлового обмена.

Контекст момента

Apache Kafka 0.9 вышла в начале 2016 с новым consumer API и поддержкой Kerberos-аутентификации

На одном из проектов по интеграции у нас несколько лет живёт классическая схема: ERP ночью выгружает транзакции в CSV, скрипт подтягивает файл, парсит, льёт в аналитическое хранилище. Работает. Но «работает» здесь значит «никто не жалуется утром» - а не «всё хорошо».

Проблем три. Первая - данные в хранилище отстают на ночь. Аналитики получают свежую картину только после обработки пакета, а любой отчёт в 10 утра - это вчерашняя реальность. Вторая - файловый обмен хрупкий: упала сетевая папка, изменилась кодировка в выгрузке, добавилось поле - всё это обнаруживается утром с задержкой. Третья - при добавлении нового потребителя данных (новый отчёт, новая система) схема не масштабируется изящно: либо плодишь отдельные выгрузки из ERP, либо гоняешь один и тот же файл в несколько мест.

В начале года вышла Kafka 0.9, и мы решили посмотреть, насколько она решает эти три проблемы в нашем конкретном случае.

Что изменилось в 0.9

Kafka существует несколько лет, но 0.9 - заметный релиз. Два момента, которые нас интересовали.

Новый consumer API. Старый high-level consumer API работал через ZooKeeper для отслеживания офсетов - что создавало дополнительную зависимость и периодически приводило к неожиданному rebalance. В 0.9 офсеты хранятся в специальном топике внутри самой Kafka. Consumer group API стал чище и предсказуемее.

Kerberos/SASL. Kafka до 0.9 жила в предположении, что брокеры находятся в доверенной сети. Для корпоративной инфраструктуры с AD и Kerberos это был серьёзный аргумент против. Теперь SASL/GSSAPI поддерживается нативно - можно не городить отдельный VPN-туннель между ERP-хостом и Kafka-брокерами.

Как выглядит пилот

Нагрузка на продакшн-ERP у клиента - порядка 5000 транзакций в секунду в пиковые часы. Это не огромно по меркам Kafka, но для нас это первый реальный поток, а не синтетический тест.

Схема пилота:

ERP (producer) -> Kafka topic "transactions" -> consumer -> DWH (PostgreSQL)

На стороне ERP написан небольшой producer: каждая зафиксированная транзакция публикуется в топик в JSON. Топик с replication factor 2, три брокера на выделенных VM. Consumer читает батчами, трансформирует и вставляет в хранилище через COPY.

Первое, что проверяли - гарантии доставки. Kafka даёт at-least-once из коробки: сообщение считается доставленным после подтверждения от брокера (параметр acks). Consumer сам управляет офсетом и коммитит его только после успешной записи в DWH. Дубликаты возможны при падении consumer между записью и коммитом офсета - это надо учитывать на стороне хранилища. У нас транзакции имеют уникальный ID, так что идемпотентная вставка через UPSERT закрывает этот случай. Это же, кстати, то, что мы уже делали с PostgreSQL 9.5 INSERT ON CONFLICT.

Что получается на практике

По задержке - данные в хранилище появляются с отставанием в несколько секунд вместо ночи. Это уже другой режим работы аналитики.

По надёжности - брокер честно держит сообщения на диске. Когда в ходе пилота мы уронили consumer на 40 минут (не специально - деплой пошёл не так), после перезапуска он прочитал отставший батч и догнал поток. ERP ничего не заметил, данные в DWH не потерялись.

По масштабируемости - добавить второго consumer для нового потребителя данных означает новый consumer group с отдельным офсетом. ERP при этом не трогаем вообще. Это и есть pub/sub в чистом виде.

Из неожиданного: настройка ZooKeeper-кластера (Kafka по-прежнему требует ZooKeeper для координации брокеров) заняла больше времени, чем сам Kafka. Три ноды ZooKeeper, нечётное число для кворума - это не сложно, но надо помнить, что это отдельный компонент с отдельными требованиями к latency между нодами.

Мониторинг оказался отдельным разговором. Kafka экспортирует метрики через JMX. Подключить к нашему Grafana+InfluxDB стеку через JMX-коллектор - решаемо, но с танцами. Смотрим на consumer lag как на основной индикатор здоровья потока: если consumer не успевает за producer, lag растёт - это сигнал, что что-то пошло не так.

Где пока неудобно

Kafka 0.9 - не коробочное решение. Здесь нет веб-интерфейса из коробки, нет встроенного мониторинга с дашбордами, нет простого способа посмотреть «что вообще лежит в топике» без консольных утилит. Для разработки и отладки это ощутимо.

Schema для сообщений мы держим договорённостью: producer и consumer договорились о структуре JSON, и это пока работает. При одной паре producer-consumer - терпимо. При большем числе участников отсутствие schema registry начнёт болеть, но это задача следующего этапа, если пилот покажет результат.

Операционно: Kafka - это Java, это JVM, это heap. На трёх брокерах отдали по 6 ГБ heap каждому. Следить за GC паузами и размером heap надо отдельно.

Итог первого месяца

Пилот работает, данные идут. Ночной ETL-файл пока не отключили - гоним параллельно и сверяем. Расхождений нет, что само по себе хороший знак.

Задача следующего этапа - убедить заказчика переключить аналитику на real-time поток как основной источник и убрать ночную выгрузку как костыль. Это уже не техническая задача, а процессная.

Контакт

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

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