A B C D E F G H I J K L M N O P Q R S T V W Y Z А Б В Г Е И К М О П С Т Ц

Change Data Capture

Change Data Capture

 

Change Data Capture (CDC) это подход к отслеживанию изменений данных в базе и передаче событий об операциях INSERT, UPDATE и DELETE в другие системы почти в реальном времени. Технология относится к классу решений для интеграции и репликации данных: вместо регулярной полной выгрузки таблиц приёмник получает только то, что действительно изменилось. Типовая цепочка выглядит так: PostgreSQL отдаёт изменения, коннектор Debezium превращает их в события, Apache Kafka хранит поток, ClickHouse принимает его в аналитическое хранилище.

 

Что такое Change Data Capture и какую задачу он закрывает

Классическая интеграция данных долго строилась на пакетной выгрузке. Ночью запускается job, вычитывает таблицу целиком или по условию и грузит её в хранилище. Схема работает годами, но у неё два врождённых дефекта. Данные в приёмнике отстают на весь интервал между запусками, а нагрузка на источник растёт пропорционально размеру таблицы, а не количеству правок.

Захват изменений решает обе проблемы разом. Система фиксирует сам факт изменения строки в момент коммита транзакции и отдаёт его наружу отдельным сообщением. Объём трафика теперь зависит только от интенсивности записи. Если за час в таблице на 500 миллионов строк поменялось три сотни записей, по проводу уедут ровно эти три сотни событий.

Важно понимать, что CDC это не продукт, а именно паттерн. Реализаций много: коннекторы Debezium, встроенная функциональность СУБД, облачные сервисы репликации, оркестраторы потоков. Механика у них похожая, отличается способ добычи изменений и гарантии доставки.

 

Архитектура CDC-пайплайна

Практически любой поток захвата изменений раскладывается на четыре звена, и понимание их границ сильно помогает при отладке.

  • Источник. Транзакционная база, где живут исходные данные: PostgreSQL, MySQL, Oracle, SQL Server, MongoDB. Именно она ведёт журнал транзакций, из которого всё вычитывается.
  • Коннектор. Компонент, который читает журнал, декодирует бинарный формат и превращает каждую правку в структурированное событие с ключом, состоянием до и состоянием после. В экосистеме Kafka эту роль обычно играет Debezium внутри Kafka Connect.
  • Транспорт. Брокер сообщений, который хранит поток событий и развязывает источник с потребителями. Чаще всего Apache Kafka, реже Pulsar или Kinesis.
  • Приёмник. Аналитическая СУБД, озеро данных, поисковый индекс или сервис, которому нужны свежие данные. Здесь события применяются к целевой модели.

Такое разделение даёт запас прочности. Если приёмник упал, события никуда не пропадают, они лежат в топике брокера и будут прочитаны после восстановления.

Change Data Capture, путь изменения строки от журнала WAL PostgreSQL через слот репликации и Debezium в Kafka и аналитическое хранилище

 

 

Apache Kafka для инженеров данных

Код курса
DEVKI
Ближайшая дата курса
24 августа, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800

 

Как работает log-based CDC под капотом

Любая надёжная СУБД пишет изменения сначала в журнал, и только потом в файлы данных. В PostgreSQL это журнал упреждающей записи (write-ahead log, WAL), в MySQL двоичный журнал (binary log, binlog), в Oracle redo log. Журнал нужен базе для восстановления после сбоя, но у него есть побочное свойство: там уже лежит полная и упорядоченная история всех правок. Журнальный захват изменений просто подключается к этому потоку вторым читателем. Ничего лишнего база при этом не делает, журнал пишется в любом случае.

 

WAL, логическое декодирование и слоты репликации

Сырой WAL это физический формат, привязанный к номерам страниц на диске, читать его напрямую бесполезно. Поэтому PostgreSQL предлагает механизм логического декодирования (logical decoding). Он проходит по журналу и разворачивает физические записи в понятные строковые события с именами таблиц и значениями колонок. Формат вывода задаёт плагин, начиная с десятой версии в комплекте идёт штатный pgoutput, а для отладки удобен test_decoding с человекочитаемым текстом.

Второй ключевой элемент это слот репликации (replication slot). Слот запоминает позицию, до которой потребитель успешно прочитал журнал, и не даёт серверу удалить более свежие сегменты WAL. Благодаря слоту коннектор переживает перезапуск и продолжает ровно с той точки, где остановился, ничего не потеряв. Механизм выключен по умолчанию: пока параметр wal_level стоит в значении replica, данных строк в журнале просто нет.

Логическое декодирование работает только с закоммиченными транзакциями, так что откаченные изменения в поток не попадают. Порядок событий соответствует порядку коммитов, а не порядку самих операций внутри транзакций.

 

Три способа поймать изменение и почему выигрывает журнал

Журнальный подход не единственный. До него в ходу были ещё два, и они до сих пор встречаются там, где доступа к журналу нет.

Критерий Опрос по метке времени Триггеры в базе Чтение журнала
Как ловит правку Периодический SELECT с условием по колонке updated_at Триггер пишет копию строки в служебную таблицу Внешний процесс читает WAL или binlog
Ловит DELETE Нет, удалённой строки уже нет в выборке Да Да
Нагрузка на источник Средняя, регулярные сканы по индексу Высокая, каждая запись дорожает Низкая, журнал пишется в любом случае
Задержка Равна интервалу опроса, обычно минуты Секунды Доли секунды
Изменения схемы приложения Нужна колонка с меткой времени Нужен DDL на каждую таблицу Не нужны
Промежуточные состояния строки Теряются между опросами Сохраняются Сохраняются

Опрос по метке времени проще всего внедрить, но он молча теряет удаления и промежуточные версии строки. Триггеры честнее, зато утяжеляют каждую транзакцию и превращают схему в лес служебных объектов. Журнальный вариант выигрывает по всем важным параметрам, и цена за это одна: нужен доступ к репликационному интерфейсу базы, а его дают не всегда.

 

Подводные камни, о которых редко пишут в туториалах

CDC выглядит элегантно на схеме и заметно менее элегантно в эксплуатации. Вот что стреляет чаще всего. Список короткий, но злой.

  • Зависший слот репликации. Слот держит WAL до тех пор, пока потребитель не подтвердит чтение. Остановили коннектор на выходные и забыли, а журнал всё это время копится и однажды забивает диск источника. Мониторинг отставания слота обязателен, это не опция.
  • Первичный снимок. Журнал хранит только недавнюю историю, поэтому перед стримингом коннектор делает initial snapshot и вычитывает таблицы целиком. На большой базе это долгая и тяжёлая операция, которую надо планировать отдельно.
  • Эволюция схемы. Разработчик добавил колонку, и формат событий поменялся. Без реестра схем (schema registry) и договорённостей о совместимости потребители ломаются в тот же момент.
  • Дубликаты при перезапуске. Большинство реализаций дают гарантию at-least-once, а не exactly-once. Приёмник должен быть идемпотентным, иначе после сбоя в аналитике появятся задвоенные строки.
  • Операции, которые журнал не описывает построчно. TRUNCATE приезжает отдельным событием без данных строк, а DDL в потоке выглядит пустой транзакцией. И то и другое требует отдельной обработки на стороне приёмника.

Change Data Capture, почему непрочитанный слот репликации удерживает сегменты WAL и заполняет диск базы-источника

Ни один из этих пунктов не отменяет пользу подхода, но каждый превращается в инцидент, если про него узнать уже в проде.

 

 

Практическая архитектура данных

Код курса
PRAR
Ближайшая дата курса
14 сентября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800

 

Когда CDC нужен, а когда это перебор

Сценарии применения узнаваемы. Их примерно пять. Потоковая аналитика, где витрины в ClickHouse или другом OLAP-хранилище должны отставать от оперативной базы на секунды, а не на сутки. Репликация между разнородными системами, когда штатная физическая реплика не годится из-за разных СУБД. Наполнение поискового индекса и кэшей, чтобы они не расходились с источником. Разбор монолита на сервисы, когда новые сервисы читают события старой базы. Аудит изменений, где важна полная история версий строки.

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

 

Логическое декодирование PostgreSQL 18 на прогоне

Чтобы увидеть механизм своими глазами, Debezium и Kafka не нужны. Всё, что делает коннектор, доступно средствами самой базы. Ниже реальный прогон на PostgreSQL 18.4 с параметром wal_level в значении logical и плагином test_decoding.

По традиции весь код используемый в статье выкладываем на наш GitHub репозиторий

Change Data Capture (CDC) from airflow import DAG from airflow.operators.bash import BashOperator from datetime import datetime with DAG( dag_id="spark_submit_demo", start_date=datetime(2025, 1, 1), schedule="@daily", catchup=False ) as dag: run = BashOperator( task_id="run_job", bash_command="spark-submit app.py" ) GitHub code example Change Data Capture (CDC)

Файл cdc_slot_demo.sql создаёт таблицу заказов, заводит слот репликации и прогоняет три операции в трёх отдельных транзакциях.

-- PostgreSQL 18.4 (Homebrew), плагин логического декодирования test_decoding
-- прогон на стенде 2026-08-14, полный вывод в run_output.txt

-- шаг 1: убираем следы прошлого прогона, скрипт должен запускаться повторно
SELECT pg_drop_replication_slot(slot_name) FROM pg_replication_slots
 WHERE slot_name = 'cdc_demo_slot';
DROP TABLE IF EXISTS orders;

-- шаг 2: таблица-источник, за изменениями которой мы следим
CREATE TABLE orders (
    id       bigint PRIMARY KEY,
    customer text NOT NULL,
    status   text NOT NULL,
    amount   numeric(10,2) NOT NULL
);

-- шаг 3: слот репликации, он же точка чтения журнала и стоп-кран для его очистки
SELECT pg_create_logical_replication_slot('cdc_demo_slot', 'test_decoding');

-- шаг 4: обычная работа приложения, три операции в трёх транзакциях
INSERT INTO orders VALUES (1, 'ООО Ромашка', 'new', 1500.00);
UPDATE orders SET status = 'paid' WHERE id = 1;
DELETE FROM orders WHERE id = 1;

-- шаг 5: вычитываем поток изменений и сдвигаем позицию слота вперёд
SELECT lsn, xid, data FROM pg_logical_slot_get_changes('cdc_demo_slot', NULL, NULL);

Последний запрос возвращает девять строк. Каждая транзакция обрамлена парой BEGIN и COMMIT, а между ними лежит ровно одно событие. Вот колонка data этого вывода.

CDC на Postrgresql

 

Полнота строки в UPDATE и DELETE

Обратите внимание на последнее событие. От удалённого заказа приехал только идентификатор, ни имени клиента, ни суммы в потоке нет. Это поведение по умолчанию: база кладёт в журнал лишь первичный ключ старой версии строки, потому что этого достаточно для физической репликации. Приёмнику, который строит историю изменений, такого события мало.

Лечится это режимом REPLICA IDENTITY FULL на конкретной таблице. После него в событии UPDATE появляется блок old-key с прежними значениями всех колонок, а DELETE везёт строку целиком. Цена известна заранее: журнал растёт, потому что в него теперь пишется вдвое больше данных.

table public.orders: UPDATE: old-key: id[bigint]:2 customer[text]:'АО Василёк' status[text]:'new' amount[numeric]:990.00 new-tuple: id[bigint]:2 customer[text]:'АО Василёк' status[text]:'paid' amount[numeric]:1090.00
table public.orders: DELETE: id[bigint]:2 customer[text]:'АО Василёк' status[text]:'paid' amount[numeric]:1090.00

 

Цена молчащего потребителя

Самый частый инцидент эксплуатации проверяется тем же слотом. Файл cdc_slot_lag.sql вставляет двести тысяч строк, пока потребитель молчит, и замеряет, сколько журнала база вынуждена держать на диске.

-- PostgreSQL 18.4 (Homebrew), отставание слота репликации, прогон на стенде 2026-08-14
-- запускается после cdc_slot_demo.sql

-- шаг 1: имитируем рабочую нагрузку, пока потребитель молчит
INSERT INTO orders SELECT g, 'клиент ' || g, 'new', g * 1.5
  FROM generate_series(3, 200000) g;

-- шаг 2: ключевая метрика эксплуатации, сколько журнала база держит из-за слота
SELECT slot_name, active,
       pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), confirmed_flush_lsn)) AS retained_wal
  FROM pg_replication_slots;

-- шаг 3: потребитель прочитал поток, слот сдвинулся, журнал освободился
SELECT count(*) AS events_read FROM pg_logical_slot_get_changes('cdc_demo_slot', NULL, NULL);
SELECT slot_name,
       pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), confirmed_flush_lsn)) AS retained_wal
  FROM pg_replication_slots;

На прогоне 199998 вставленных строк удержали 31 МБ журнала. Потребитель вычитал 200 тысяч событий, слот сдвинулся, и метрика вернулась к нулю. Масштабируйте цифру на рабочую базу с потоком в тысячи транзакций в секунду и коннектором, который лежит вторые сутки, и вы получите точный сценарий заполнения диска. Именно поэтому запрос из шага 2 стоит завести в мониторинг раньше, чем сам пайплайн уедет в прод.

 

От слота к Debezium и Kafka

Debezium делает ровно то же самое, только вместо ручного запроса подключается к слоту по протоколу репликации и раскладывает события по топикам Kafka. В Kafka Connect коннектор описывается обычным JSON, где указаны координаты базы, имя слота и публикации.

// синтаксис по документации Debezium 3.6, PostgreSQL connector
{
  "name": "orders-cdc-connector",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "postgres",
    "database.port": "5432",
    "database.user": "debezium",
    "database.dbname": "shop",
    // логический префикс топиков, события уедут в shop_cdc.public.orders
    "topic.prefix": "shop_cdc",
    "table.include.list": "public.orders",
    // штатный плагин логического декодирования, значение по умолчанию
    "plugin.name": "pgoutput",
    "slot.name": "dbz_orders_slot",
    "publication.name": "dbz_publication",
    // initial: сначала полный снимок таблицы, затем стриминг из WAL
    "snapshot.mode": "initial"
  }
}

Дальше начинается обычная инженерия потоков. Как эти события обрабатываются на лету и что делать с гарантиями доставки, разбирается на курсе Практическая архитектура данных.

 

Заключение

Change Data Capture переносит интеграцию данных из режима периодических выгрузок в режим непрерывного потока событий. Технически это чтение журнала транзакций, который база и так ведёт для собственного восстановления, поэтому источник почти не страдает. Настройка коннектора занимает пару файлов конфигурации, и это самая лёгкая часть работы. Основное усилие лежит в эксплуатации: мониторинг отставания слота, план первичного снимка, договорённости об эволюции схемы и идемпотентный приёмник. Разобравшись с этими четырьмя пунктами, вы получите поток изменений, на который можно опереться в проде.

 

Референсные ссылки