Change Data Capture (CDC): Postgres WAL → Debezium → Kafka topics → fan-out (Elasticsearch search index, Snowflake DWH, PaymentService, NotificationService). Demonstrates logical replication slots with pgoutput plugin, initial snapshot + streaming switchover, transactional outbox pattern for atomic business events, replication slot growth disaster (Debezium down → WAL накапливается → диск Postgres переполняется), and schema evolution via ALTER TABLE. Includes ADRs on CDC vs dual-write vs sync API call (ADR-001) and replication slot growth multi-layer protection (ADR-002).
Change Data Capture (CDC) — паттерн, при котором приложение пишет только в свою OLTP базу (Postgres / MySQL / Mongo), а отдельный процесс (Debezium) читает её WAL/binlog и эмитит каждое изменение row как event в Kafka. Downstream-системы (поиск, кэш, DWH, другие микросервисы) подписываются и реагируют. Источник правды один — журнал, который СУБД и так пишет для durability.
Альтернатива — dual-write (app в одной операции бьёт в DB и в Kafka) — даёт консистентность ровно до первого failure. Через год накапливаются сотни разошедшихся записей, и никто не помнит почему поиск показывает старое.
App делает INSERT/UPDATE/COMMIT в Postgres — единственная атомарная операция. Postgres всё равно пишет change в WAL (это обязательно для crash recovery). Logical replication slot — отдельная сущность Postgres — хранит позицию (LSN), до которой Debezium подтвердил приём. Слот ГАРАНТИРУЕТ не удалить WAL до этой позиции, чтобы consumer мог восстановиться после reconnect. Debezium через streaming-replication-протокол pullит изменения, декодирует pgoutput-bytes → logical row event (op=c/u/d, before/after, LSN, txId), сериализует в Avro и публикует в Kafka topic shop.public.<table>. Партиционирование по PK даёт strict ordering per row. Latency от COMMIT до event — single-digit ms. Downstream фанаутятся параллельно: Elasticsearch re-index, Snowflake MERGE, микросервис реагирует на OrderPlaced.
При первом запуске connector'а в downstream нет ничего. Debezium делает initial snapshot: открывает транзакцию REPEATABLE READ, фиксирует start-LSN, читает SELECT * FROM table батчами по 10k (lock-free, через snapshot isolation) и эмитит каждый row как op=r (read) с флагом snapshot=true. После окончания snapshot'а переключается на streaming с зафиксированного LSN — zero data loss, ни одного event между snapshot и streaming не пропускается. Для огромных таблиц (1B+ rows) есть incremental snapshot (signal-based, KIP-211) — можно re-snapshot диапазон без stop'а connector'а или снять snapshot с read-replica, чтобы не блокировать prod.
Raw CDC эмитит физические row diffs — для аналитики и репликации это идеально, но микросервису-подписчику нужен бизнес-event с правильной семантикой (OrderPlaced{order_id, customer, total}), а не «UPDATE orders SET status='placed' WHERE id=42». Решение: outbox table. В той же транзакции app пишет в основную таблицу И в outbox(event_type, payload). Атомарность гарантирует Postgres. Debezium с SMT (Single Message Transform) роутит outbox table в отдельный topic orders.events. Получаешь явный event API + atomicity + чистую границу между бизнес-логикой и сырыми row changes. Cleanup outbox через cron DELETE WHERE created_at < NOW() - 7d.
Главный operational риск. Если Debezium встал (OOM, K8s evict, Kafka недоступна, broker rebalance), confirmed_flush_lsn замораживается → Postgres не ротирует WAL → pg_wal/ растёт со скоростью write workload (100 MB/min = 144 GB/сутки). Когда диск заполняется, Postgres переходит в read-only, потом крашится — outage основной OLTP базы из-за CDC-инструмента, который должен был быть «невидимой репликацией». Multi-layer защита: (1) max_slot_wal_keep_size = 50GB — hard cap, после которого PG отъединяет slot и удаляет WAL (Debezium при reconnect получит slot invalidated и сделает re-snapshot — больно, но не outage), (2) Prometheus alert pg_current_wal_lsn() - confirmed_flush_lsn > 5 GB → on-call paged, (3) auto-restart Debezium через Kafka Connect REST endpoint.
ALTER TABLE orders ADD COLUMN phone — Debezium читает new relation OID из WAL, re-fetch'ит schema, регистрирует новую Avro schema v2 в Schema Registry с BACKWARD compatibility (nullable field). Старые consumer'ы (ES sink) игнорируют unknown field — не падают. Новые consumer'ы (Snowflake) видят новую колонку. Без Schema Registry schema changes ломают consumer'ов silently — на масштабе это deadly.
INSERT/UPDATE/COMMIT в Postgresshop.public.orders (raw CDC), shop.public.outbox (outbox table CDC), orders.events (бизнес-events после SMT-роутинга)ADR-001 на OrderService разбирает выбор CDC vs dual-write vs sync API call. ADR-002 на replication slot — multi-layer защиту от slot growth disaster.
MERGECDC log-based vs Dual-write vs Sync API. Dual-write (INSERT в DB + Producer.send в один блок кода) НЕ atomic — на 100k tx/day по эмпирике 1–10 потерь/сутки, через год тысячи разъездов. Sync API (POST в каждый downstream) складывает latency (50→200ms), создаёт n×n coupling, любой downstream down → весь OrderService падает. CDC log-based: один write в DB, latency 1–10ms до event, at-least-once + per-row ordering, single source of truth = WAL. Цена — operational complexity Debezium + Kafka Connect + Schema Registry.
Raw CDC vs Outbox. Raw CDC = физические row diffs, идеален для репликации/аналитики/search. Outbox = логические бизнес-events с правильной семантикой, для микросервисной интеграции. Best practice: оба паттерна одновременно — raw для analytics, outbox для events.
Log-based vs Trigger-based vs Polling. Triggers (DB-triggers пишут в audit table) — portable, но кладут DB load на каждую запись + ordering issues при concurrency. Polling (SELECT WHERE updated_at > $last) — простой, но high latency (≥ period), full table scans, не видит DELETE, пропускает изменения с одинаковым timestamp. Log-based (WAL/binlog) — gold standard для high-throughput: minimal DB load (WAL пишется всё равно), strict ordering, низкая latency. Cons — setup complexity + replication slot growth.
pgoutput vs wal2json. Встроенный pgoutput (PG10+) — binary, эффективный, не требует extension. wal2json — JSON output, человекочитаемый, удобен для отладки, но deprecated для production CDC. Используй pgoutput.
At-least-once vs Exactly-once. Debezium гарантирует at-least-once (при reconnect возможны дубли). Consumer ОБЯЗАН дедуплицировать по (LSN, txId) или (table, pk, op). Иллюзия exactly-once без consumer dedup — путь к двойным charge на PaymentService.
binlog_format=STATEMENT для MySQL. Debezium не получит row-level changes. Только ROW + binlog_row_image=FULLtable.include.list. Иначе bloat: каждая миграция, каждая audit-таблица улетают в Kafkaop=d + before (+ tombstone для compaction). Если consumer не handle — данные на downstream остаются призраками