ClickHouse умеет читать сообщения напрямую из Kafka через Kafka Engine.
Для этого настраивается цепочка из трёх таблиц и одного представления.
📥 Kafka-таблица
Описывает структуру сообщений и Kafka-топик, из которого будут читаться данные.
CREATE TABLE default.kafka_orders
(
`id` Int32,
`status` String,
`price` String,
`__deleted` Nullable(String)
)
ENGINE = Kafka('broker:9092', 'inventory.orders', 'clickhouse', 'AvroConfluent')
SETTINGS format_avro_schema_registry_url = 'http://schema-registry:8081';
🔁Материализатор данных из Kafka
Kafka-таблица читает сообщения только один раз — смещения коммитаются в consumer group.
Поэтому каждую запись нужно сразу перекладывать в постоянную таблицу.
CREATE MATERIALIZED VIEW default.consumer__orders
TO default.stream_orders
(
`id` Int32,
`status` String,
`price` String,
`__deleted` Nullable(String)
) AS
SELECT
id,
status,
price,
__deleted
FROM default.kafka_orders;
🧱 Основная таблица
Хранит все версии строк и пометки об удалении.
Для корректной замены старых записей используется ReplacingMergeTree.
CREATE TABLE default.stream_orders
(
`id` Int32,
`status` String,
`price` String,
`__deleted` String
)
ENGINE = ReplacingMergeTree
ORDER BY (id, price)
SETTINGS index_granularity = 8192;
👀 Витрина данных
Скрывает удалённые строки и возвращает только актуальное состояние данных.
CREATE VIEW default.orders
(
`id` Int32,
`status` String,
`price` String
) AS
SELECT
id,
status,
price
FROM default.stream_orders
FINAL
WHERE __deleted = 'false';
Важно:
постоянное использование FINAL дорого по ресурсам.
В production лучше:
1. Агрегации, - last value
2. Фоновые merge, - ожидание схлопывания данных
3. Материализованные витрины, - предрасчитанные представления
✅ Заключение
Мы собрали полноценный конвейер синхронизации между MySQL и ClickHouse через CDC.
Ключевые элементы:
1. Debezium, - читает binlog MySQL
2. Kafka, - гарантирует доставку и порядок событий
3. Kafka Engine, - потоковая загрузка в ClickHouse
4. ReplacingMergeTree, - устранение дубликатов
5. Поле __deleted, - корректная обработка удалений
В результате получается аналитическая копия боевой OLTP-базы:
MySQL продолжает обслуживать транзакции,ClickHouse — тяжёлую аналитику и отчёты,
оба без взаимных блокировок и деградации производительности.
На этом мы завершаем линейку постов про ClickHouse.Мы разобрали, на мой взгляд, все ключевые аспекты — дальше только практика-практика и еще раз практика
Ещё услышимся 👋
#ClickHouse #MySQL #CDC #Debezium #Kafka #DataSync #OLTP #OLAP #DataEngineering #АналитикаДанных