Содержание
- Что такое Debezium и какую задачу он закрывает
- Архитектура и способы развёртывания
- Debezium на Kafka Connect
- Debezium Server и встраиваемый Engine
- Как работает захват изменений под капотом
- Снимок и потоковая вычитка
- Структура события изменения
- Три вещи, которые видно только на прогоне
- Чем CDC по журналу лучше опроса таблиц
- Ограничения и подводные камни
- Когда применять Debezium, а когда не стоит
- Практика, минимальный стенд без Kafka
- Заключение
- Референсные ссылки
Debezium (Debezium) это open-source платформа захвата изменений данных (change data capture, CDC), которая читает журнал транзакций базы данных и превращает построчные изменения в поток событий. Платформа отслеживает операции INSERT, UPDATE и DELETE в PostgreSQL, MySQL, MariaDB, SQL Server, Oracle, Db2, MongoDB и ряде других СУБД. Чаще всего Debezium работает в связке с Apache Kafka и Kafka Connect, доставляя изменения в аналитические витрины, Data Lake и микросервисы почти в реальном времени.
Что такое Debezium и какую задачу он закрывает
Debezium относится к классу log-based CDC-инструментов. Ключевое слово здесь log-based: платформа не опрашивает таблицы запросами, а подписывается на тот же журнал, по которому база восстанавливается после сбоя и кормит свои реплики. В PostgreSQL это write-ahead log через логическое декодирование, в MySQL и MariaDB это binlog, в SQL Server отдельные CDC-таблицы, в Oracle LogMiner или XStream. Общий механизм разобран в отдельной статье про Change Data Capture, здесь речь про конкретную реализацию.
Задача звучит просто, но решается тяжело. Данные лежат в операционной СУБД, а нужны они в другом месте: в хранилище для аналитики, в поисковом индексе, в кэше, в соседнем сервисе. Ночной batch-выгрузкой эту потребность уже не закрыть, бизнесу нужна свежесть в секундах. Debezium даёт поток, где каждое изменение строки становится отдельным сообщением.
Важный нюанс. Debezium не хранилище и не движок обработки. Это транспорт, точнее его первое звено. Платформа отвечает за то, чтобы ни одно изменение не потерялось по дороге из журнала в брокер сообщений. Что происходит дальше, решают потребители.
Архитектура и способы развёртывания
Архитектурно Debezium это набор коннекторов, по одному на СУБД, плюс общее ядро с логикой снимка, потоковой вычитки и хранения смещений. Коннекторы разделяют общий дизайн, поэтому знание одного переносится на остальные почти без потерь. Развернуть набор можно тремя способами, и выбор влияет на всю остальную инфраструктуру.
Debezium на Kafka Connect
Основной и самый распространённый вариант. Коннектор ставится плагином в кластер Kafka Connect, а тот берёт на себя скучную, но критичную работу: распределение задач по воркерам, перезапуск после падения, хранение смещений и истории схемы в служебных топиках. Отказоустойчивый рантайм не надо писать самому, поэтому связка и стала фактическим стандартом. Разбор такого контура на живом примере есть в статье CDC-репликация Big Data в реальном времени с Apache Kafka и Debezium. Если Kafka Connect пока выглядит чёрным ящиком, стоит начать с курса Apache Kafka для инженеров данных, где рантайм разбирается руками.
Debezium Server и встраиваемый Engine
Не всем нужна Kafka. Для таких случаев есть два обходных пути.
- Debezium Server. Standalone-приложение, которое читает изменения из СУБД и отправляет их в другую инфраструктуру: Pulsar, Kinesis, Google Pub/Sub, Redis Streams, RabbitMQ, NATS, обычный HTTP-эндпоинт. Kafka Connect при этом не нужен вообще.
- Debezium Engine. Библиотека, которая встраивается прямо в Java-приложение, события обрабатываются в собственном коде. Максимум контроля и минимум внешних зависимостей, но отказоустойчивость и масштабирование теперь ваша забота.
Правило выбора грубое, но рабочее. Есть Kafka в контуре, берите Kafka Connect. Kafka нет и не планируется, смотрите на Debezium Server. Engine оставьте на случай, когда нужна встроенная логика, которую неудобно выносить в отдельный сервис.
Apache Kafka для инженеров данных
Код курса
DEVKI
Ближайшая дата курса
24 августа, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800
Как работает захват изменений под капотом
Механика делится на две фазы, и путать их не стоит. Первая отвечает на вопрос, что уже есть в базе. Вторая на вопрос, что в ней меняется прямо сейчас.
Снимок и потоковая вычитка
При первом запуске коннектор снимает начальный снимок (snapshot): вычитывает содержимое отслеживаемых таблиц и порождает по событию на каждую существующую строку. Так потребитель получает полную картину, а не обрывок с середины. Закончив снимок, коннектор переключается в потоковый режим (streaming) и дальше читает только журнал.
Поведением снимка управляет параметр snapshot.mode. Значений несколько: initial снимает картину один раз, always повторяет её при каждом старте, no_data пропускает данные и берёт только схему, initial_only снимает и останавливается. Отдельного упоминания заслуживает инкрементальный снимок (incremental snapshot). Он позволяет добрать данные по таблице на живом коннекторе, не останавливая поток и не переснимая всё заново. Штука незаменимая, когда в отслеживаемый набор понадобилось добавить таблицу спустя полгода после запуска.
Структура события изменения
Событие устроено как конверт: состояние строки до операции, состояние после, метаданные источника и код операции. Ниже реальное событие удаления, снятое на стенде с Debezium Server 3.6.1.Final и PostgreSQL 18.
# событие DELETE, прогон на Debezium Server 3.6.1.Final и PostgreSQL 18
{
"before": {"id": 2, "customer": "", "status": "", "amount": "AA=="},
"after": null,
"source": {
"version": "3.6.1.Final",
"connector": "postgresql",
"name": "shop",
"db": "shopdb",
"schema": "public",
"table": "orders",
"txId": 771,
"lsn": 29375448,
"snapshot": "false"
},
"op": "d",
"ts_ms": 1786744488208
}
Поле op кодирует тип операции: c для вставки, u для обновления, d для удаления, r для строки, прочитанной во время снимка, t для очистки таблицы. Секция source это паспорт события: имя коннектора, база, схема, таблица, идентификатор транзакции и позиция в журнале. По позиции всегда понятно, откуда пришло изменение и в каком порядке оно случилось.
Три вещи, которые видно только на прогоне
Документация описывает счастливый путь, а стенд показывает детали, на которых спотыкаются первые же интеграции. Прогон на связке PostgreSQL 18 и Debezium Server 3.6.1.Final дал три штуки разом.
- Блок before по умолчанию пустой. В событии выше в нём заполнен только идентификатор, остальные поля отдают пустышками. В событии UPDATE весь блок вообще равен null. Причина в REPLICA IDENTITY таблицы: журнал по умолчанию везёт только первичный ключ, и взять старое значение попросту неоткуда.
- Числовые типы приезжают в base64. Сумма 1500.00 в событии выглядит как строка AA== или ATiA, потому что numeric по умолчанию кодируется двоичным Decimal. Лечится параметром decimal.handling.mode со значением string или double.
- Публикация создаётся на все таблицы. Список table.include.list фильтрует события уже внутри коннектора, а публикация в PostgreSQL по умолчанию создаётся для всех таблиц базы. Декодируется при этом лишний трафик, а лечится параметром publication.autocreate.mode со значением filtered.
Первый пункт чинится на стороне базы командой ALTER TABLE orders REPLICA IDENTITY FULL. После неё событие обновления приезжает с заполненным before, и дельта наконец становится полноценной.
# то же событие UPDATE после REPLICA IDENTITY FULL, прогон на стенде
{
"before": {"id": 4, "customer": "umbrella", "status": "shipped", "amount": "ATiA"},
"after": {"id": 4, "customer": "umbrella", "status": "closed", "amount": "ATiA"},
"op": "u",
"ts_ms": 1786744587706
}
Цена решения честная: журнал начинает везти полную старую версию строки, объём WAL растёт. Включать REPLICA IDENTITY FULL стоит там, где дельта действительно нужна, а не на всех таблицах подряд.
Чем CDC по журналу лучше опроса таблиц
Прежде чем ставить платформу, полезно понять, от чего вы уходите. Захват изменений исторически решался тремя способами, и у каждого своя цена.
| Критерий | Log-based CDC (Debezium) | Опрос по updated_at | Триггеры в БД |
|---|---|---|---|
| Нагрузка на исходную БД | Низкая, читается журнал | Растёт с частотой опроса | Высокая, на каждой записи |
| Ловит DELETE | Да | Нет, строка исчезает бесследно | Да |
| Промежуточные состояния | Видны все | Теряются между опросами | Видны все |
| Изменение схемы приложения | Не требуется | Нужна колонка со временем | Нужны объекты в БД |
| Задержка | Секунды и меньше | Равна интервалу опроса | Секунды |
| Сложность эксплуатации | Высокая | Низкая | Средняя |
Главный проигрыш опроса не в задержке, а в потерях. Строку удалили, и запрос по updated_at про это никогда не узнает. Строка успела трижды сменить статус между опросами, и вы увидите только последний. Триггеры эти проблемы решают, но платят производительностью транзакций и превращают схему базы в место, куда заглядывают с опаской.
Построение DWH на ClickHouse
Код курса
CLICH
Ближайшая дата курса
12 октября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800
Ограничения и подводные камни
Log-based CDC не бесплатный обед. Проблемы у него другие, чем у опроса, но они есть, и знать про них лучше заранее.
- Слот репликации в PostgreSQL. Логическое декодирование держит слот, а слот держит WAL. Коннектор упал и не подтверждает позицию, WAL растёт, диск на боевой базе кончается. Мониторинг лага слота обязателен, это самая частая авария в проде.
- Права и настройка источника. Нужен пользователь с правами репликации, параметр wal_level со значением logical в PostgreSQL, ROW-формат binlog в MySQL, включённый CDC на уровне базы и таблиц в SQL Server. Прийти к администратору базы придётся в любом случае.
- Изменения схемы. Платформа отслеживает DDL и ведёт историю схемы, но потребители к смене типа колонки или удалению поля автоматически не готовы. Без Schema Registry и договорённости о совместимости такое изменение ломает приёмник.
- Семантика at-least-once. Дубли событий возможны при перезапуске после сбоя. Потребитель обязан быть идемпотентным, обычно за счёт ключа и upsert-логики на приёмнике.
- Объём потока. Одно изменение строки даёт одно сообщение с полным конвертом. На таблице с массовыми обновлениями поток раздувается быстро, а политика хранения и компакции топиков становится частью проектирования.
Ни один пункт не повод отказаться от платформы. Но каждый превращается в инцидент, если про него не подумали заранее.
Когда применять Debezium, а когда не стоит
Платформа хорошо ложится на вполне конкретный класс задач.
- Репликация в аналитический контур. Наполнение DWH, Data Lake или колоночного хранилища свежими данными без ночного batch-окна.
- Синхронизация производных хранилищ. Поисковый индекс, кэш, материализованные представления в другой системе, которые должны догонять источник за секунды.
- Развязка сервисов и паттерн outbox. Сервис пишет событие в таблицу-аутбокс внутри своей транзакции, платформа вычитывает её и публикует в брокер. Классическое решение проблемы двойной записи.
- Постепенная миграция с legacy-системы. Старая база остаётся источником истины, новая догоняет её через поток изменений.
А вот где инструмент избыточен. Если данные обновляются раз в сутки и никого не смущает утренняя выгрузка, CDC добавит кластер и дежурство по слотам, не дав ничего взамен. Если нужна не дельта, а агрегат по всей таблице, проще посчитать его запросом. Если источник это внешний API, а не СУБД с журналом транзакций, читать нечего и подход неприменим.
Практика, минимальный стенд без Kafka
По традиции весь код используемый в статье выкладываем на наш GitHub репозиторий
Самый быстрый способ пощупать платформу руками это Debezium Server, которому нужен один контейнер вместо кластера. Файл docker-compose.yml поднимает PostgreSQL 18 с включённым логическим декодированием и сам сервер.
# docker-compose.yml, прогнан на Docker 29.4.2, образ quay.io/debezium/server:3.6.1.Final
services:
postgres:
image: postgres:18
container_name: dbz_pg
environment:
POSTGRES_PASSWORD: postgres
POSTGRES_DB: shopdb
# логическое декодирование включается параметрами запуска сервера
command:
- postgres
- -c
- wal_level=logical
- -c
- max_replication_slots=10
- -c
- max_wal_senders=10
ports:
- "5433:5432"
volumes:
- ./init.sql:/docker-entrypoint-initdb.d/init.sql:ro
server:
image: quay.io/debezium/server:3.6.1.Final
container_name: dbz_server
depends_on:
- postgres
volumes:
- ./config/application.properties:/debezium/config/application.properties:ro
- ./data:/debezium/data
Вся настройка коннектора живёт в файле application.properties. Он же показывает минимальный набор параметров, без которых сервер не стартует.
# application.properties, прогнан на Debezium Server 3.6.1.Final debezium.sink.type=http debezium.sink.http.url=http://host.docker.internal:8099/events debezium.source.connector.class=io.debezium.connector.postgresql.PostgresConnector debezium.source.offset.storage.file.filename=/debezium/data/offsets.dat debezium.source.offset.flush.interval.ms=1000 # координаты исходной базы debezium.source.database.hostname=postgres debezium.source.database.port=5432 debezium.source.database.user=postgres debezium.source.database.password=postgres debezium.source.database.dbname=shopdb # логический префикс, отслеживаемые таблицы и имя слота репликации debezium.source.topic.prefix=shop debezium.source.table.include.list=public.orders debezium.source.plugin.name=pgoutput debezium.source.slot.name=debezium_shop debezium.source.snapshot.mode=initial # формат событий, без этих двух строк сервер падает на старте debezium.format.key=json debezium.format.value=json debezium.format.value.schemas.enable=false debezium.format.key.schemas.enable=false
Параметр topic.prefix задаёт пространство имён, table.include.list сужает захват до нужных таблиц, slot.name называет тот самый слот, за лагом которого придётся следить. Пара строк про формат выглядит необязательной, но без них сервер выпадает с ошибкой SRCFG00014 прямо на старте.
Конверт с полями before и after удобен для аудита и неудобен для приёмника, которому нужна просто строка. Для этого случая есть преобразование ExtractNewRecordState, оно же unwrap. С ним то же событие удаления выглядит совсем иначе.
# вывод после SMT unwrap и decimal.handling.mode=string, прогон на стенде
{"id": 4, "customer": "umbrella", "status": "new", "amount": "777.25",
"__deleted": "false", "__op": "c", "__table": "orders", "__lsn": 29391336}
{"id": 1, "customer": "acme", "status": "new", "amount": "1500.00",
"__deleted": "true", "__op": "d", "__table": "orders", "__lsn": 29391704}
Плоская строка вместо конверта, сумма нормальным числом, тип операции в служебном поле. Именно в таком виде события обычно и уезжают в приёмник, а полный конверт остаётся там, где нужна история изменений.
Заключение
Debezium превращает журнал транзакций СУБД в упорядоченный поток событий и делает это, не трогая ни схему исходной базы, ни путь записи в неё. Взамен он требует внимания к эксплуатации: слоты репликации, REPLICA IDENTITY, эволюция схемы и идемпотентность потребителей перестают быть чужой заботой. Если поток изменений действительно нужен в секундах, цена оправдана. Если хватает суточной выгрузки, честнее обойтись без CDC.



