Skip to content
Learning Platform

Ментальная модель: база уже ведёт стрим изменений за вас

Главная идея, которую нужно усвоить до всего остального: реляционная база данных уже является стриминговой системой. Она не «хранит таблицы» в том наивном смысле, в каком это рисуют в учебниках. Прежде чем хоть один байт страницы будет переписан на диске, СУБД записывает намерение изменить данные в последовательный журнал — WAL в Postgres, binlog в MySQL, redo log в Oracle. Этот журнал — упорядоченный, durable поток фактов «что произошло с данными». Именно из него реплики догоняют мастер, и именно его читает crash recovery после падения.

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 решает это через logical decoding — слой, который прогоняет WAL через output-плагин (pgoutput, встроенный с Postgres 10) и выдаёт уже понятные события: в таблице public.orders вставлена строка с такими значениями.

Чтобы это заработало, нужны три вещи на стороне базы:

  • wal_level = logical в конфиге — Postgres начинает писать в WAL дополнительную информацию, достаточную для логического восстановления строк, а не только физических страниц.
  • replication slot — серверный курсор, отмечающий, до какого места Debezium уже прочитал журнал.
  • Publication — набор таблиц, которые попадают в поток.

Replication slot — самая опасная и самая важная деталь. Слот гарантирует, что Postgres не удалит сегменты WAL, пока их не прочитал потребитель. Это и есть фундамент надёжности: упал Debezium на час — журнал его дождётся. Но обратная сторона: если коннектор умер навсегда, а слот остался, WAL растёт без ограничений и в какой-то момент забивает диск мастера. Заброшенный слот — классическая причина ночного инцидента. Мониторинг отставания слота (pg_replication_slots.confirmed_flush_lsn против текущего LSN) обязателен.

Позиция в журнале адресуется через LSN. Это монотонно растущий оффсет; Debezium хранит его как свой курсор и при рестарте просит базу «продолжай с этого 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 — не самостоятельный сервис, а набор плагинов для Kafka Connect. Это снимает с него огромный класс задач: распределение работы по воркерам, рестарты при падении, REST API для управления и, главное, durable-хранение оффсетов в служебных топиках Kafka.

Поток данных линеен: база пишет 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 — сигнал для механизма log compaction в Kafka: при компакции топика этот ключ будет физически удалён, и топик останется корректным представлением текущего состояния таблицы. Если sink игнорирует 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 даёт at-least-once, и это не недоработка, а честное следствие физики системы. Представьте: Debezium прочитал изменение из WAL, отправил его в Kafka, Kafka подтвердил запись — и в этот момент, до того как Debezium успел закоммитить новый оффсет (LSN) в свой служебный топик, процесс падает. После рестарта коннектор продолжит с последнего закоммиченного оффсета — то есть переотправит уже доставленное событие. Дубликат.

Сделать 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.

Ещё в направлении · Data Engineering

Все материалы направления →