Почти любая современная СУБД ведёт журнал изменений, куда сначала записываются все операции, а уже потом они применяются к данным.
Это механизм Write Ahead Log 📝
В MySQL таким журналом является binlog.
Если читать этот журнал, интерпретировать изменения и передавать их в другую систему, мы фактически реализуем подход Change Data Capture (CDC) 🔄
Почему CDC удобен для синхронизации:
- Работает потоково, почти без задержек
- Обеспечивает конечную согласованность
- Не требует тяжёлых батч-процедур
- Сохраняет порядок изменений
- Масштабируется под высокую нагрузку
Для работы с binlog чаще всего используют Debezium. Он подключается к MySQL, отслеживает изменения и публикует их в Kafka через Kafka Connect.
Дальше эти события может забирать ClickHouse 🧩
Ниже — только те настройки Debezium, которые важны именно для корректной работы с ClickHouse.
1. Оставляем только актуальное состояние записи
По умолчанию Debezium отправляет событие в таком виде:
- Cостояние до изменения
- Cостояние после изменения
Плюс при удалении формируется «пустое» сообщение
Для Kafka это нормально, а вот для ClickHouse — неудобно: таблицы Kafka Engine ожидают плоскую структуру.
Поэтому включаем трансформацию:
"transforms": "unwrap",
"transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState"
Что это даёт:
- Для INSERT и UPDATE остаётся только новое состояние строки
- Поле before отбрасывается
- Структура сообщения упрощается
2. Корректно обрабатываем удаления
После упрощения структуры Debezium перестаёт передавать удаления. Чтобы это исправить, добавляем:
"transforms.unwrap.delete.handling.mode": "rewrite"
Теперь:
- Удалённые записи не пропадают
- В сообщении появляется поле __deleted = true
- Остальные операции получают __deleted = false
Это позволяет:
- Хранить историю изменений
- Фильтровать удалённые строки уже на стороне ClickHouse (через view)
3. Проблема обновлений неключевых полей
В нашем примере:
- В MySQL первичный ключ — id
- В ClickHouse таблица отсортирована по (id, status)
Если обновляется поле status, ClickHouse видит новую комбинацию ключей и создаёт ещё одну строку → появляются дубликаты 😬
Чтобы этого избежать, Debezium нужно явно сказать, какие поля считать идентификатором записи при обновлениях.
Настройка:
"message.key.columns": "inventory.orders:id;inventory.orders:status"
Теперь при изменении этих полей:
- Сначала генерируется событие удаления старой версии
- Затем событие вставки новой
- ClickHouse корректно «заменяет» строку через ReplacingMergeTree
Итоговая конфигурация Debezium
{
"name": "mysql-connector",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"database.hostname": "mysql",
"database.port": "3306",
"database.user": "root",
"database.password": "mypassword",
"database.server.id": "2",
"database.server.name": "dbz.inventory.v2",
"database.include.list": "inventory",
"table.include.list": "inventory.orders",
"message.key.columns": "inventory.orders:id;inventory.orders:status",
"schema.history.internal.kafka.bootstrap.servers": "broker:9092",
"schema.history.internal.kafka.topic": "dbz.inventory.history.v2",
"snapshot.mode": "schema_only",
"topic.prefix": "dbz.inventory.v2",
"transforms": "unwrap",
"transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
"transforms.unwrap.delete.handling.mode": "rewrite"
}
}⚠️ Важный момент: как выбирать message.key.columns
После задания message.key.columns:
- Эти поля используются как ключ сообщения в Kafka
- По ним распределяются данные по партициям
- Нарушенный порядок событий = риск рассинхронизации в ClickHouse
Практическое правило:
1️⃣ Определите ключ сортировки таблицы в ClickHouse
2️⃣ Поймите, из каких колонок источника он формируется
3️⃣ Объедините все эти колонки
4️⃣ Укажите их в message.key.columns
5️⃣ Убедитесь, что они входят в ORDER BY в ClickHouse
#ClickHouse #MySQL #CDC #Debezium #Kafka #DataSync #OLTP #OLAP #DataEngineering #АналитикаДанных