| Title | 🛠️ Debezium CDC | |||
|---|---|---|---|---|
| Group | Dev & Tools | |||
| Icon | 🛠️ | |||
| Order | 15 | |||
| tags |
|
Debezium is an open-source Change Data Capture (CDC) platform that streams database changes in near real-time. It monitors database transaction logs (WAL/binlog/oplog/redologs) and publishes changes as events, most commonly into Apache Kafka.
Common use cases / Типичные сценарии: Real-time replication, event-driven architectures, database audit trails, synchronization with search engines (OpenSearch/Elasticsearch), data lake ingestion, cache invalidation, microservices integration.
Status / Статус: Actively maintained and widely used in production. Alternatives: Apache NiFi (flow-based data integration), Maxwell's Daemon (lightweight MySQL CDC), Oracle GoldenGate (enterprise CDC), AWS DMS (managed CDC), Kafka Connect JDBC (polling-based sync).
Default ports / Порты по умолчанию: Kafka
9092, Kafka Connect REST API8083, Schema Registry8081
- Architecture
- Installation & Configuration
- Core Management
- Sysadmin Operations
- Security
- Backup & Restore
- Troubleshooting & Tools
- Production Runbooks
- Logrotate Configuration
- Additional Notes
- Official Documentation
Database → Transaction Logs (WAL/binlog/oplog) → Debezium Connector → Kafka Connect → Kafka Topics → Consumers / Search / Analytics
| Database | Transaction Log | Notes / Примечания |
|---|---|---|
| PostgreSQL | WAL | Requires logical replication / Нужна логическая репликация |
| MySQL | binlog | Requires ROW binlog format / Нужен формат ROW |
| MariaDB | binlog | Similar to MySQL / Аналогично MySQL |
| MongoDB | oplog | Replica set required / Нужен replica set |
| Oracle | redo logs | XStream/LogMiner |
| SQL Server | CDC tables | SQL Server CDC feature |
| Db2 | Transaction logs | Enterprise usage |
| Layer | Description EN | Описание RU | Best Use Case |
|---|---|---|---|
| Layer 4 | TCP-level balancing | Балансировка TCP уровня | Raw Kafka traffic |
| Layer 7 | HTTP-aware balancing | HTTP-aware балансировка | Kafka Connect REST API |
| Type | Description EN | Описание RU | Best For |
|---|---|---|---|
| Active | Load balancer probes service | Балансировщик проверяет сервис | HA production clusters |
| Passive | Detects failures from traffic | Ошибки из трафика | Simpler environments |
| Service / Сервис | Port / Порт | Description / Описание |
|---|---|---|
| Kafka Broker | 9092 |
Kafka plaintext |
| Kafka SSL | 9093 |
Kafka TLS |
| Kafka Connect REST API | 8083 |
Connect API |
| Schema Registry | 8081 |
Avro schema registry |
| PostgreSQL | 5432 |
PostgreSQL |
| MySQL | 3306 |
MySQL |
/opt/debezium/docker-compose.yml
version: '3.9'
services:
zookeeper:
image: confluentinc/cp-zookeeper:7.6.0
environment:
ZOOKEEPER_CLIENT_PORT: 2181
kafka:
image: confluentinc/cp-kafka:7.6.0
ports:
- "9092:9092"
environment:
KAFKA_BROKER_ID: 1
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://<HOST>:9092
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 # DEV ONLY: Use 3 for production
> [!WARNING]
> The `KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1` setting is for single-node development only. Production environments must run at least 3 brokers and set topic replication factors to 3.
connect:
image: debezium/connect:2.6
ports:
- "8083:8083"
environment:
BOOTSTRAP_SERVERS: kafka:9092
GROUP_ID: 1
CONFIG_STORAGE_TOPIC: connect_configs
OFFSET_STORAGE_TOPIC: connect_offsets
STATUS_STORAGE_TOPIC: connect_statusesdocker compose up -d # Start stack / Запустить стек
docker compose ps # Check containers / Проверить контейнерыSample output:
NAME STATUS
kafka Up
connect Up
zookeeper Up
/var/lib/pgsql/data/postgresql.conf
wal_level = logical # Required for CDC / Обязательно для CDC
max_replication_slots = 10 # Slots for connectors / Слоты для коннекторов
max_wal_senders = 10 # WAL sender processes / Процессы WAL sender/var/lib/pgsql/data/pg_hba.conf
host replication <USER> <IP>/32 md5 # Allow replication / Разрешить репликациюsystemctl restart postgresql # Restart PostgreSQL / Перезапустить PostgreSQLpsql -U postgresCREATE ROLE <USER> WITH REPLICATION LOGIN PASSWORD '<PASSWORD>';/etc/my.cnf
server-id=1
log_bin=mysql-bin
binlog_format=ROW # Required for CDC / Обязательно для CDC
binlog_row_image=FULL # Full row images / Полные образы строк
expire_logs_days=7 # Retention / Хранение/opt/debezium/connectors/postgres.json
{
"name": "postgres-connector",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "<HOST>",
"database.port": "5432",
"database.user": "<USER>",
"database.password": "<PASSWORD>",
"database.dbname": "appdb",
"database.server.name": "app-postgres",
"topic.prefix": "app",
"plugin.name": "pgoutput"
}
}# Register connector
curl -X POST http://<HOST>:8083/connectors \
-H "Content-Type: application/json" \
--data @postgres.jsoncurl http://<HOST>:8083/connectors/postgres-connector/status # Check status / Проверить статусSample output:
{
"name": "postgres-connector",
"connector": {
"state": "RUNNING"
}
}kafka-topics.sh \
--bootstrap-server <HOST>:9092 \
--list # List all topics / Все топикиkafka-console-consumer.sh \
--bootstrap-server <HOST>:9092 \
--topic app.public.users \
--from-beginning # Read from start / Читать сначалаSample output:
{
"before": null,
"after": {
"id": 1,
"name": "Ivan"
},
"op": "c"
}curl -X POST http://<HOST>:8083/connectors \
-H "Content-Type: application/json" \
--data @connector.json # Create / Создатьcurl http://<HOST>:8083/connectors # List all / Список всехcurl http://<HOST>:8083/connectors/<CONNECTOR>/config # Get config / Получить конфигcurl -X PUT http://<HOST>:8083/connectors/<CONNECTOR>/pause # Pause / Паузаcurl -X PUT http://<HOST>:8083/connectors/<CONNECTOR>/resume # Resume / ВозобновитьWarning
Deleting connector offsets may cause full resnapshot or duplicate events. Удаление офсетов коннектора может вызвать полный ресnapshot или дубликаты.
curl -X DELETE http://<HOST>:8083/connectors/<CONNECTOR> # Delete / Удалитьdocker ps # List containers / Список контейнеров
docker logs -f connect # Follow logs / Смотреть логи
docker restart connect # Restart service / Перезапускsystemctl status kafka-connect # Check service / Проверить сервис
systemctl restart kafka-connect # Restart service / Перезапустить
journalctl -u kafka-connect -f # Follow logs / Логи| Path / Путь | Description / Описание |
|---|---|
/var/log/kafka/server.log |
Kafka broker logs / Логи брокера |
/var/log/kafka-connect/connect.log |
Kafka Connect logs / Логи Connect |
/var/lib/kafka/data/ |
Kafka topic storage / Хранилище топиков |
/var/lib/postgresql/data/pg_wal/ |
PostgreSQL WAL files / WAL файлы |
/var/lib/mysql/ |
MySQL binlogs / Бинлоги MySQL |
/etc/kafka/connect-distributed.properties
KAFKA_HEAP_OPTS="-Xms2G -Xmx2G" # JVM heap size / Размер кучи JVM| RAM | Heap Recommendation / Рекомендация |
|---|---|
| 4 GB | 1-2 GB |
| 8 GB | 2-4 GB |
| 16 GB+ | 4-8 GB |
kafka-consumer-groups.sh \
--bootstrap-server <HOST>:9092 \
--describe \
--group connect-cluster # Check lag / Проверить лагss -lntp # Listening TCP ports / TCP порты
nc -zv <HOST> 9092 # Test Kafka port / Проверить Kafka
curl http://<HOST>:8083/ # Test Connect API / Проверить APIfirewall-cmd --permanent --add-port=9092/tcp # Allow Kafka / Разрешить Kafka
firewall-cmd --permanent --add-port=8083/tcp # Allow Connect / Разрешить Connect
firewall-cmd --reload # Apply rules / Применитьiptables -A INPUT -p tcp --dport 9092 -j ACCEPT # Allow Kafka
iptables -A INPUT -p tcp --dport 8083 -j ACCEPT # Allow Connect/etc/kafka/connect-distributed.properties
security.protocol=SSL
ssl.truststore.location=/etc/kafka/secrets/kafka.truststore.jks
ssl.truststore.password=<PASSWORD>
ssl.keystore.location=/etc/kafka/secrets/kafka.keystore.jks
ssl.keystore.password=<PASSWORD>Caution
Storing ssl.truststore.password and ssl.keystore.password in plaintext is unsafe in production. Replace inline passwords with provider-based secret lookups using Kafka Connect secrets management (config.providers), HashiCorp Vault, or Kubernetes Secrets.
keytool -genkeypair \
-alias kafka-connect \
-keyalg RSA \
-keystore kafka.keystore.jks # Generate keystore / Создать keystoreopenssl s_client -connect <HOST>:9093 # Verify TLS / Проверить TLSGRANT CONNECT ON DATABASE appdb TO <USER>;
GRANT USAGE ON SCHEMA public TO <USER>;
GRANT SELECT ON ALL TABLES IN SCHEMA public TO <USER>;GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT
ON *.* TO '<USER>'@'%';mkdir -p /opt/backups/debezium
curl http://<HOST>:8083/connectors/<CONNECTOR>/config \
-o /opt/backups/debezium/<CONNECTOR>.json # Save config / Сохранить конфигCaution
Using tar on /var/lib/kafka/data while Kafka is running can produce inconsistent snapshots. Filesystem backups should only be used when Kafka is stopped or quiesced. For production, rely on Kafka replication, use MirrorMaker2 for cross-cluster/topic replication, or use Kafka-aware backup tools that support consistent snapshots.
# ONLY WHEN KAFKA IS STOPPED
tar czf kafka-data-backup.tar.gz /var/lib/kafka/data/ # Backup Kafka data / Бэкап данных Kafkapg_dump -U <USER> appdb > appdb.sql # Dump database / Дамп БДpsql -U <USER> appdb < appdb.sql # Restore database / Восстановить БДconnect-mirror-maker.sh mm2.properties # Start MirrorMaker2 / Запустить MirrorMaker2| Problem / Проблема | Cause / Причина | Fix / Решение |
|---|---|---|
| Connector FAILED | Invalid config / Невалидный конфиг | Validate JSON / Проверить JSON |
| No events in Kafka | Replication disabled / Репликация отключена | Enable WAL/binlog |
| Duplicate events | Offset reset / Сброс офсетов | Check offsets / Проверить офсеты |
| Lag increasing | Slow consumers / Медленные консьюмеры | Scale consumers / Масштабировать |
| WAL growing | Connector stopped / Коннектор остановлен | Resume connector / Возобновить |
curl http://<HOST>:8083/connectors/<CONNECTOR>/status | jq # Check status / Проверить статусkcat -b <HOST>:9092 -L # List metadata / Метаданныеpsql -U postgres -c "SELECT * FROM pg_replication_slots;" # List slots / Список слотовmysql -e "SHOW MASTER STATUS;" # Binlog status / Статус бинлогаdu -sh /var/lib/postgresql/data/pg_wal/ # WAL size / Размер WALcurl -X POST \
http://<HOST>:8083/connectors/<CONNECTOR>/tasks/0/restart # Restart task / Перезапускjcmd <PID> VM.flags # JVM flags / Флаги JVM
jstat -gc <PID> 5s # GC stats / Статистика GC
jmap -heap <PID> # Heap info / Информация о куче- Validate database replication settings / Проверить настройки репликации
- Create dedicated replication user / Создать пользователя репликации
- Verify Kafka connectivity / Проверить связь с Kafka
- Register connector JSON / Зарегистрировать JSON коннектора
- Verify connector status / Проверить статус
- Consume test events / Прочитать тестовые события
- Configure monitoring and alerting / Настроить мониторинг
- Configure backups / Настроить бэкапы
Warning
Incorrect rollback may cause duplicate or missing events. Некорректный откат может вызвать дубликаты или потерю событий.
- Pause connector / Поставить на паузу
- Export current config / Экспортировать текущий конфиг
- Restore previous connector config / Восстановить предыдущий конфиг
- Resume connector / Возобновить работу
- Validate offsets / Проверить офсеты
- Verify event ordering / Проверить порядок событий
- Check Docker/systemd status / Проверить статус
- Review logs / Просмотреть логи
- Validate Kafka broker availability / Проверить доступность Kafka
- Validate DB connectivity / Проверить связь с БД
- Restart failed task / Перезапустить задачу
- Resume connector / Возобновить коннектор
- Verify topic ingestion / Проверить поступление данных
Caution
Full WAL/binlog storage can stop database writes. Заполненное хранилище WAL/binlog может остановить запись в БД.
- Verify connector status / Проверить статус коннектора
- Resume failed connector / Возобновить коннектор
- Increase storage / Увеличить хранилище
- Remove obsolete logs only if replication confirmed / Удалять только при подтверждённой репликации
- Validate replication lag / Проверить лаг репликации
/etc/logrotate.d/kafka-connect
/var/log/kafka-connect/*.log {
daily
rotate 14
compress
missingok
notifempty
copytruncate
}| Topic / Топик | Purpose / Назначение |
|---|---|
connect_configs |
Connector configurations / Конфигурации коннекторов |
connect_offsets |
Connector offsets / Офсеты коннекторов |
connect_statuses |
Connector states / Состояния коннекторов |
| Mode / Режим | Description EN | Описание RU |
|---|---|---|
initial |
Full initial snapshot | Полный начальный snapshot |
schema_only |
Only schema | Только схема |
never |
No snapshot | Без snapshot |
when_needed |
Snapshot if required | Snapshot при необходимости |
| Code / Код | Meaning / Значение |
|---|---|
c |
Create / Создание |
u |
Update / Обновление |
d |
Delete / Удаление |
r |
Snapshot read / Чтение snapshot |
- Use dedicated replication users / Используйте выделенных пользователей репликации
- Enable TLS in production / Включите TLS в продакшене
- Monitor replication lag / Мониторьте лаг репликации
- Store connector configs in Git / Храните конфиги коннекторов в Git
- Avoid resetting offsets in production / Не сбрасывайте офсеты в продакшене
- Use Schema Registry for Avro/Protobuf / Используйте Schema Registry
- Separate Kafka disks from OS disks / Разделяйте диски Kafka и ОС
- Configure alerting for connector failures / Настройте алертинг
- Debezium Official: https://debezium.io/documentation/
- Apache Kafka: https://kafka.apache.org/documentation/
- Kafka Connect: https://docs.confluent.io/platform/current/connect/index.html
- PostgreSQL Logical Replication: https://www.postgresql.org/docs/current/logical-replication.html
- MySQL Binary Log: https://dev.mysql.com/doc/refman/en/binary-log.html
- OpenSearch: https://opensearch.org/docs/