Debezium CDC вместо ночного ETL: изменения из PostgreSQL в Kafka за секунды
Заменили ночные ETL-выгрузки на Debezium 0.9 CDC через logical replication slot: изменения в PostgreSQL появляются в Kafka-топике за секунды, аналитические витрины обновляются в near-real-time.
Debezium 0.9 - CDC-коннектор для PostgreSQL через logical replication slot в Kafka Connect
Когда мы в мае разбирали KSQL для real-time агрегации POS-событий, там источником уже была Kafka - события писались туда напрямую с кассовых терминалов. Этот пост про другую ситуацию: транзакционная система на PostgreSQL, трогать её нельзя, а аналитические витрины нужно кормить не раз в сутки, а непрерывно. Классическая задача, и решение через Change Data Capture наконец стало достаточно зрелым, чтобы тащить в продакшн без ощущения что идёшь по минному полю.
Откуда задача
Клиент - сервисная компания с основной БД на PostgreSQL. Там живёт всё: клиенты, договоры, заявки, история состояний. DWH на отдельном кластере, аналитические витрины строятся для отчётности и оперативного контроля.
До нас схема была типичной: ночной ETL забирал дельту по полю updated_at, трансформировал, грузил в аналитику. Работало. Проблема - «ночной». К утру витрины отставали на сутки, и для части задач это стало неприемлемо: менеджеры хотели видеть состояние активных заявок в реальном времени, а не по состоянию на начало дня.
Казалось бы - запускай ETL чаще. Запустили каждые 15 минут. Получили нагрузку на OLTP-базу: каждые 15 минут тяжёлый SELECT по updated_at с большим диапазоном, конкуренция с рабочими транзакциями, периодические таймауты. Откатились. Начали думать.
Почему CDC, а не polling
Polling через updated_at - дешёвый способ получить дельту, но с известными ограничениями. Первое - удалённые строки не поймать: строка удалена, updated_at нет, в выгрузку не попадёт. Второе - нагрузка на OLTP пропорциональна частоте опроса. Третье - транзакционная атомарность не гарантирована: если между двумя запусками ETL транзакция успела зафиксироваться и ещё одна изменила те же строки, в выгрузку попадёт промежуточное состояние или второе изменение перезапишет первое незамеченным.
CDC работает принципиально иначе: читает журнал транзакций. Для PostgreSQL это WAL, через механизм logical decoding. Событие появляется в потоке ровно один раз, в порядке коммита, с полной информацией: какая строка изменилась, что было до и что стало после. Удаление - тоже событие. Нагрузка на OLTP минимальная - читается WAL, который Postgres пишет в любом случае.
Debezium 0.9 на практике
Debezium запускается как коннектор в Kafka Connect. Мы уже использовали Kafka Connect для других коннекторов, так что дополнительная инфраструктура не понадобилась - добавили ещё один воркер.
На стороне PostgreSQL нужно три вещи. Первое - включить логическую репликацию: в postgresql.conf выставить wal_level = logical. Второе - создать publication для таблиц, которые хотим отслеживать. Третье - выдать пользователю Debezium права репликации.
-- публикация только нужных таблиц
CREATE PUBLICATION debezium_pub FOR TABLE
clients, contracts, requests, request_history;
-- пользователь с правами на репликацию
CREATE ROLE debezium REPLICATION LOGIN PASSWORD '...';
GRANT SELECT ON clients, contracts, requests, request_history TO debezium;
Конфигурация коннектора минималистичная:
{
"name": "pg-cdc-connector",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "pg-oltp.internal",
"database.port": "5432",
"database.user": "debezium",
"database.dbname": "main",
"database.server.name": "oltp",
"plugin.name": "pgoutput",
"publication.name": "debezium_pub",
"slot.name": "debezium_slot",
"table.whitelist": "public.clients,public.contracts,public.requests,public.request_history",
"transforms": "unwrap",
"transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState"
}
}
После старта коннектор делает initial snapshot - читает текущее состояние таблиц и пишет в Kafka-топики. Потом переключается на WAL и пишет только изменения. Каждая таблица - отдельный топик: oltp.public.clients, oltp.public.requests и так далее.
Что в топике
Каждое сообщение - JSON с конвертом Debezium: до (before), после (after), операция (c/u/d), временная метка транзакции, LSN в WAL. Мы используем ExtractNewRecordState transform, который разворачивает конверт и оставляет только финальное состояние строки - это упрощает консьюмеры на стороне аналитики, им не нужно знать про формат Debezium.
Для history-таблиц ExtractNewRecordState не нужен - там интересны все события включая удаления, так что читаем полный конверт.
Как это вписалось в витрины
На стороне DWH поставили Kafka Connect с JDBC Sink коннектором на аналитическую базу. Изменения из Kafka-топиков идут напрямую в таблицы-источники DWH. Витрины обновляются через триггеры и incremental materialized views - это отдельная история про ClickHouse, но суть в том, что от коммита в OLTP до обновления витрины проходит несколько секунд, а не 15 минут или ночь.
Нагрузку на PostgreSQL мы замерили до и после. Периодические всплески от ETL-запросов пропали полностью. Replication slot добавляет незначительный overhead - Postgres должен держать WAL до тех пор, пока коннектор не подтвердит чтение. Пока коннектор работает штатно, это несколько секунд WAL, ничего страшного. Важно мониторить pg_replication_slots - если коннектор упал и перестал читать, WAL начнёт накапливаться и место на диске кончится. Это реальный риск, который надо держать на алертах.
Грабли, которые нашли
Logical replication slot и failover. Replication slot привязан к конкретному серверу PostgreSQL. Если OLTP переключается на реплику (плановый failover или аварийный), slot на старом мастере умирает, на новом мастере его нет. Debezium не умеет автоматически создать slot на новом мастере и начать с правильного LSN - требуется ручное вмешательство. У клиента есть Patroni, и этот момент потребовал отдельной процедуры обслуживания при switchover.
DDL-изменения. Если в OLTP добавляют новую колонку, Debezium продолжает работать, но схема топика в Schema Registry не обновляется автоматически. Нужно явно эволюционировать схему или рестартовать коннектор. Несколько раз поймали рассинхрон, пока не настроили процедуру для DDL-изменений в OLTP.
initial snapshot и размер таблиц. Если таблица большая, initial snapshot занимает время - коннектор читает всю таблицу в рамках одной транзакции, удерживает snapshot-транзакцию. Для таблицы request_history в несколько десятков миллионов строк это заняло несколько часов. Во время снапшота WAL всё равно читается и буферизуется - после снапшота коннектор дочитывает накопившееся. Это нормальная ситуация, просто надо понимать что первый старт - не быстрый.
Промежуточный итог
Ночной ETL больше не нужен. Изменения в OLTP появляются в Kafka-топиках за секунды. Аналитические витрины живут в near-real-time. OLTP не нагружается периодическими SELECT-ами.
Debezium 0.9 - это не «поставил и забыл». Replication slots требуют внимания, DDL-изменения требуют процедуры, failover требует ручной работы. Но по сравнению с polling-ETL это принципиально другой уровень свежести данных при меньшей нагрузке на источник. По проектам с такими задачами мы работаем в рамках DWH и аналитики.