Kafka → ClickHouse с материализованным представлением: заменили двухчасовую SSIS-джобу на потоковый ETL
Убрали ночную SSIS-джобу, которая крутилась два часа, и заменили связкой Kafka + ClickHouse Kafka Engine + Materialized View. Данные доступны в реальном времени.
ClickHouse активно развивается в 2019: материализованные представления и движок Kafka позволяют строить потоковые ETL-пайплайны
У клиента была SSIS-джоба. Это такая вещь, которую в Microsoft SQL Server называют "интеграционным сервисом", а в реальности это XML-пакет весом 12 мегабайт, который два часа перекладывает данные из одной базы в другую. Каждую ночь. И за два часа успевало сломаться что-нибудь одно: то коннект к источнику, то нехватка памяти на сервере, то просто "ошибка 0xC020801C, обратитесь к администратору".
Клиент пришёл к нам с задачей для DWH и аналитики: данные о транзакциях из нескольких источников должны быть доступны аналитикам не к утру следующего дня, а - желательно - в реальном времени или хотя бы с лагом в минуты.
Что было и почему не работало
Архитектура сложилась исторически: несколько операционных баз на MSSQL, SSIS-пакет собирал данные ночью, выгружал в staging-таблицы, потом хранимые процедуры агрегировали всё в витрины. Утром аналитики открывали Power BI и видели вчерашнее.
Проблем там было несколько, и SSIS - лишь верхушка.
Batch-природа всего пайплайна. Источники работают 24/7, данные приходят непрерывно, но ETL запускался раз в сутки. Не потому что так задумывалось - просто "так исторически сложилось, и никто не переделывал".
Хрупкость. SSIS-пакет написан человеком, который уволился в 2016-м. Логика частично в пакете, частично в хранимых процедурах, частично задокументирована в голове у одного старшего аналитика. При каждой поломке час уходит на диагностику.
Масштаб растёт, джоба не. Объём данных за год вырос вдвое, время выполнения пакета с 40 минут выросло до двух часов. По этой траектории через несколько месяцев окно "ночная загрузка" просто перестанет влезать в ночь.
Почему Kafka + ClickHouse
ClickHouse мы уже использовали для событийной аналитики в другом проекте - там он хорошо показал себя на колоночных запросах с агрегациями. Здесь задача похожая: много строк, аналитические запросы с GROUP BY, фильтрами по датам и статусам.
Kafka появилась как способ разорвать жёсткую связку "источник -> ETL -> витрина". Если источники пишут события в топики, дальше их можно читать независимо, не ломая операционные базы лишними селектами в прайм-тайм.
В ClickHouse есть Kafka Engine - табличный движок, который читает сообщения из топика напрямую, без промежуточного потребителя. Выглядит это так: создаёшь таблицу с движком Kafka, указываешь брокеры и топик, и ClickHouse сам становится консьюмером. Поверх неё создаёшь материализованное представление, которое при каждой порции данных трансформирует и вставляет результат в целевую таблицу с MergeTree.
CREATE TABLE transactions_queue (
ts DateTime,
account_id UInt64,
amount Decimal(18, 2),
status LowCardinality(String)
) ENGINE = Kafka
SETTINGS
kafka_broker_list = 'kafka:9092',
kafka_topic_list = 'transactions',
kafka_group_name = 'clickhouse-consumer',
kafka_format = 'JSONEachRow';
CREATE MATERIALIZED VIEW transactions_mv TO transactions AS
SELECT
ts,
account_id,
amount,
status
FROM transactions_queue;
Материализованное представление здесь не кеш агрегата - это триггер на вставку. Данные попадают из Kafka в очередь, представление их забирает и кладёт в MergeTree. Агрегатные витрины строятся уже поверх MergeTree обычными запросами.
Как переезжали
Операционные источники не трогали сразу - там продакшн, и лезть в схему MSSQL без необходимости никто не хотел. Вместо этого поставили Debezium: он читает transaction log MSSQL и публикует изменения в Kafka как CDC-события. Источники ничего не знают о том, что их читают.
Первые две недели оба пайплайна работали параллельно: старый SSIS и новый Kafka → ClickHouse. Сверяли агрегаты - расхождения были, пришлось разбираться с тем как SSIS обрабатывал несколько edge cases в логике статусов транзакций. Оказалось, в хранимой процедуре был хардкод, который тихо игнорировал определённый тип записей. В новом пайплайне это всплыло сразу, потому что данные шли через явную трансформацию в представлении, а не через XML-пакет.
После двух недель сверки выключили SSIS. Аналитики не заметили - данные стали появляться быстрее, а не медленнее.
Что получилось на практике
Лаг данных упал с "следующее утро" до нескольких минут - это время от события в операционной базе до видимости в ClickHouse. Большую часть этих минут занимает CDC: Debezium читает лог не мгновенно.
ClickHouse на той же машине справляется с нагрузкой аналитических запросов заметно лучше, чем витрины на MSSQL справлялись раньше. Запросы за месяц с группировкой по нескольким измерениям - секунды вместо минут.
Операционная нагрузка на MSSQL-источники не выросла: CDC читает лог, а не делает SELECT-ы.
Единственное что требует аккуратности - мониторинг консьюмер-группы в Kafka. Если ClickHouse по какой-то причине отстал от топика, нужно знать об этом сразу, а не когда аналитик спросит почему данные за сегодня неполные. Настроили алерт на consumer lag в Prometheus - пока срабатывал один раз при плановом рестарте ClickHouse, и то это был ложный алерт.
Где сейчас
Пайплайн работает примерно полтора месяца. Серьёзных инцидентов не было - что уже контрастирует с историей SSIS-джобы, которая ломалась раз в две-три недели. Kafka Engine в ClickHouse производит впечатление рабочей вещи, а не эксперимента - хотя объём у нас небольшой, и делать из него далеко идущие выводы рано.
Следующий вопрос - агрегатные материализованные представления непосредственно в ClickHouse: AggregatingMergeTree позволяет держать агрегаты инкрементально без пересчёта. Смотрим, стоит ли переносить часть витринной логики туда или оставить агрегации на стороне BI-инструмента.