Ментальная модель: база уже ведёт стрим изменений за вас
Главная идея, которую нужно усвоить до всего остального: реляционная база данных уже является стриминговой системой. Она не «хранит таблицы» в том наивном смысле, в каком это рисуют в учебниках. Прежде чем хоть один байт страницы будет переписан на диске, СУБД записывает намерение изменить данные в последовательный журнал —
Change Data Capture (CDC) — это переиспользование того же журнала для внешних потребителей. Вместо того чтобы спрашивать базу «что изменилось?», вы подключаетесь к её внутреннему потоку изменений и читаете его так же, как это делает реплика. База даже не знает, что слушатель — не ещё одна реплика, а Debezium.
Если запомнить только одно: CDC не опрашивает базу, он подписывается на её журнал репликации. Всё остальное в этой статье — следствия этой одной мысли.
Почему не batch-polling
Альтернатива, которую все пробуют первой, — периодически опрашивать таблицу запросом вида SELECT * FROM orders WHERE updated_at > :last_seen. Это batch-polling, и он ломается на проде по предсказуемым причинам.
| Свойство | Batch-polling (SELECT ... WHERE updated_at >) | CDC из журнала (Debezium) |
|---|---|---|
| Нагрузка на источник | Полные сканы или индекс по updated_at, конкурируют с OLTP | Чтение журнала, который пишется в любом случае |
| Латентность | Равна интервалу опроса (секунды–минуты) | Десятки–сотни миллисекунд |
DELETE | Невидимы — удалённой строки нет в результате | Видны как отдельное событие |
| Промежуточные изменения | Теряются: две правки между опросами схлопываются в одну | Каждое изменение — отдельное событие |
| Требование к схеме | Обязателен надёжный монотонный updated_at | Не нужен |
| Транзакционные границы | Размываются | Сохраняются (txId, LSN) |
Ключевая боль — DELETE и промежуточные состояния. Опрос видит только текущий снимок строки. Если запись успели обновить дважды и удалить между двумя опросами, polling не узнает об этом ничего. Журнал же фиксирует каждый переход как отдельную запись, в строгом порядке. Для аудита, репликации в озеро данных и event-driven архитектур это разница между «примерно знаем итог» и «знаем точную историю».
Logical decoding: как Postgres отдаёт изменения наружу
WAL в исходном виде — это физический журнал: «на странице 42481 байты с такого-то смещения стали такими». Внешнему потребителю эти страничные дельты бесполезны, ему нужны логические изменения строк. Postgres решает это через pgoutput, встроенный с Postgres 10) и выдаёт уже понятные события: в таблице public.orders вставлена строка с такими значениями.
Чтобы это заработало, нужны три вещи на стороне базы:
wal_level = logicalв конфиге — Postgres начинает писать в WAL дополнительную информацию, достаточную для логического восстановления строк, а не только физических страниц.replication slotИменованный объект в Postgres, который запоминает позицию (LSN) последнего вычитанного потребителем изменения и гарантирует, что WAL до этой позиции не будет удалён. — серверный курсор, отмечающий, до какого места Debezium уже прочитал журнал.- Publication — набор таблиц, которые попадают в поток.
Replication slot — самая опасная и самая важная деталь. Слот гарантирует, что Postgres не удалит сегменты WAL, пока их не прочитал потребитель. Это и есть фундамент надёжности: упал Debezium на час — журнал его дождётся. Но обратная сторона: если коннектор умер навсегда, а слот остался, WAL растёт без ограничений и в какой-то момент забивает диск мастера. Заброшенный слот — классическая причина ночного инцидента. Мониторинг отставания слота (pg_replication_slots.confirmed_flush_lsn против текущего LSN) обязателен.
Позиция в журнале адресуется через
Стоит отдельно осознать, насколько это дёшево для источника по сравнению с polling. Чтение журнала — это последовательное чтение append-only файла, которое Postgres всё равно обязан поддерживать ради собственной репликации и рекавери. Logical decoding добавляет накладные расходы на декодирование WAL в логические записи, но не порождает конкурентных блокировок на пользовательских таблицах и не вытесняет горячие страницы из буферного кэша, как это делает полный скан при SELECT *. Для OLTP-базы под нагрузкой разница принципиальна: CDC-потребитель почти невидим для транзакционного трафика, тогда как агрессивный polling конкурирует с ним за тот же буферный пул и те же индексы.
Важно и то, что у каждой базы своя реализация того же принципа. В MySQL роль WAL играет binlog, и читать его нужно в режиме ROW (а не STATEMENT), иначе вместо значений строк в журнал попадёт текст SQL-запросов, из которых восстановить точные дельты строк нельзя. Debezium прячет эти различия за единым форматом change-event, но понимать, что под ним лежит конкретный журнал конкретной СУБД со своими ограничениями, необходимо для эксплуатации.
Debezium как source-коннектор Kafka Connect
Debezium — не самостоятельный сервис, а набор плагинов для
Поток данных линеен: база пишет WAL → logical decoding отдаёт логические изменения → Debezium source task их забирает → конвертер (обычно JSON или Avro со Schema Registry) сериализует → запись летит в топик Kafka. По умолчанию каждая таблица маппится в отдельный топик: serverName.public.orders.
Конфиг минимального Postgres-коннектора:
{
"name": "orders-connector",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "pg-primary",
"database.dbname": "shop",
"plugin.name": "pgoutput",
"slot.name": "debezium_orders",
"publication.autocreate.mode": "filtered",
"table.include.list": "public.orders",
"topic.prefix": "shop",
"snapshot.mode": "initial"
}
}
topic.prefix (он же логическое имя сервера) — это namespace во всех именах топиков и в поле source. Менять его после запуска нельзя: сменится имя топиков и поломаются оффсеты downstream.
Структура change-event: before, after, op
Каждое событие изменения — это конверт с описанием перехода строки, а не просто «новое значение». Это второй принцип, который нужно усвоить: Debezium стримит дельты состояния, а не снимки. Полезная нагрузка содержит before (строка до изменения), after (после), op (тип операции) и блок source с метаданными происхождения.
{
"op": "u",
"ts_ms": 1717200000123,
"before": { "id": 42, "status": "pending", "total": 100 },
"after": { "id": 42, "status": "paid", "total": 100 },
"source": {
"db": "shop", "schema": "public", "table": "orders",
"lsn": 23985612, "txId": 5512, "snapshot": "false"
}
}
Значения op: c (create/insert), u (update), d (delete), r (read — строка, прочитанная во время снапшота). Для c поле before равно null; для d null будет after.
Отдельно стоит понять delete. Когда строку удаляют, Debezium шлёт событие с op: "d" и after: null, а следом — tombstone: запись с тем же ключом и null-значением. Tombstone — сигнал для механизма
Блок source несёт lsn и txId — это и есть ваша связь с транзакционными границами источника. По ним downstream может группировать изменения одной транзакции и дедуплицировать.
Две фазы: snapshot и streaming
Коннектор не может начать чтение журнала «с пустого места» — в существующей таблице уже лежат миллионы строк, которых нет в текущем хвосте WAL (старые сегменты давно удалены). Поэтому жизненный цикл делится на две фазы.
Snapshot — начальная фаза. Debezium фиксирует текущий LSN, затем читает существующие строки таблицы (SELECT) и эмитит их как события с op: "r". Это «опорное» полное состояние. Современный Debezium использует incremental snapshot (алгоритм DBLog с watermark-окнами): таблица читается чанками без долгой блокирующей транзакции, и snapshot можно догнать без остановки потока.
Streaming — основная фаза. Закончив снапшот, Debezium переключается на чтение WAL начиная с зафиксированного LSN и дальше живёт в реальном времени, эмитя c/u/d. Граница спроектирована так, чтобы не было дыр: изменения, случившиеся во время снапшота, гарантированно попадут в стрим. Из-за этого нормально и ожидаемо увидеть для одной строки сначала событие r (из снапшота), а затем u (из стрима, если строку правили параллельно).
Режим снапшота настраивается параметром snapshot.mode, и выбор тут не косметический. initial (по умолчанию) делает полный снимок один раз при первом запуске, дальше — только стрим. never пропускает снапшот и читает только хвост журнала — годится, когда историческое состояние уже выгружено другим способом и нужны лишь будущие изменения. initial_only снимает снимок и останавливается, не переходя в стрим, — для разовых миграций. Понимание этих режимов критично при пересоздании коннектора: наивный рестарт с initial на большой таблице повторно прогонит весь снапшот и зальёт downstream миллионами r-событий.
Здесь же кроется тонкость с порядком. События снапшота и стрима пишутся в одни и те же топики, и для одного ключа r всегда предшествует последующим u/d. Поэтому идемпотентный приёмник, применяющий события как upsert по первичному ключу, корректно сойдётся к актуальному состоянию независимо от того, пришла строка из снапшота или из живого потока. Это ещё одна причина, по которой upsert по PK — не деталь реализации sink, а часть контракта всей CDC-цепочки.
Гарантии доставки: at-least-once и почему downstream обязан быть идемпотентным
Debezium даёт
Сделать exactly-once в общем случае невозможно дёшево: это потребовало бы атомарного коммита «запись в Kafka + сдвиг оффсета в источнике» через распределённую транзакцию, чего между Postgres и Kafka нет. Поэтому контракт честный: данные не теряются, но дубликаты возможны.
Отсюда железное правило: downstream обязан быть
- Upsert по первичному ключу. Применяйте
afterкакINSERT ... ON CONFLICT DO UPDATE. ПовторUPDATEтой же строки тем же значением — no-op. - Дедупликация по LSN/
txId. Храните последний применённыйlsnна ключ; событие со старым или равным LSN отбрасывайте. - Ключ Kafka = первичный ключ строки. Тогда все события одной строки попадают в один партишн и сохраняют порядок — без этого
uможет обогнатьc.
Порядок гарантируется только внутри партишна. Поэтому ключевание событий по PK — не оптимизация, а условие корректности: при ключевании по PK последовательность изменений одной строки никогда не переставится.
Где это даёт рычаг
CDC — это шов, который развязывает OLTP-базу и всё, что хочет реагировать на её изменения: репликация в data lake, синхронизация поисковых индексов и кэшей, выгрузка в аналитику без ночных батчей, event-driven микросервисы через паттерн outbox. Везде один и тот же фундамент — журнал репликации как источник истины и идемпотентный приёмник на другом конце.
Глубже эту механику — incremental snapshot на проде, поведение слотов под нагрузкой, схемная эволюция и outbox — разбираем на практике в бесплатном курсе Debezium CDC с песочницей: поднимаете Postgres, Kafka Connect и Debezium и сами ломаете и чините поток. Сам курс — часть бесплатного направления Data Engineering.
Что забрать с собой
Три мысли. Первая: база уже стримит изменения в журнал репликации — CDC просто читает его, а не опрашивает таблицы. Вторая: change-event описывает переход (before/after/op), а не снимок, и несёт LSN для привязки к транзакциям. Третья: гарантия — at-least-once, поэтому идемпотентный приёмник с upsert по PK и дедупликацией по LSN не опция, а обязательное условие.
Готовы потрогать руками? Начните бесплатно в песочнице курса Debezium CDC.