Ответ
Преимущества CDC:
- Минимальное воздействие на источник: CDC считывает изменения из журналов транзакций (WAL в PostgreSQL, binlog в MySQL), а не через запросы
SELECT, что практически не создаёт дополнительной нагрузки на рабочую БД. - Доставка изменений в near real-time: Позволяет строить актуальные аналитические витрины и синхронизировать системы с задержкой в секунды, что критично для систем мониторинга или персонализированных рекомендаций.
- Полный аудит изменений: Захватываются все операции (INSERT, UPDATE, DELETE), что даёт возможность восстанавливать историческое состояние данных и отслеживать, кто и что изменил.
- Надёжность: Так как логи транзакций — это механизм обеспечения целостности самой СУБД, CDC обеспечивает надежную доставку всех изменений без потерь.
Недостатки и сложности CDC:
- Сложность начальной настройки: Требует включения и настройки логирования на источнике, создания публикаций (в PostgreSQL) или назначения прав на чтение binlog.
- Обработка схемы данных: Изменения структуры таблиц (ALTER TABLE) нужно корректно обрабатывать в коннекторе (например, Debezium) и downstream-системах. Это добавляет сложности в управлении версиями.
- Зависимость от возможностей СУБД: Не все базы данных или их managed-версии в облаке поддерживают эффективный CDC. Например, могут быть ограничения на доступ к логам.
- Управление большими транзакциями: Очень большие транзакции (например, массовое обновление миллионов строк) могут создавать задержки в обработке или требовать специальной настройки коннектора.
Пример настройки Debezium для PostgreSQL:
-- На стороне базы-источника (PostgreSQL)
CREATE PUBLICATION airflow_pub FOR TABLE sales, customers;
ALTER TABLE sales REPLICA IDENTITY FULL; -- Чтобы в лог попадали старые значения строк при UPDATE/DELETE
# Конфигурация коннектора Debezium
name=inventory-connector
connector.class=io.debezium.connector.postgresql.PostgresConnector
database.hostname=source-db-host
database.dbname=mydb
database.user=replicator
database.password=***
plugin.name=pgoutput
publication.name=airflow_pub
slot.name=debezium_slot
# Трансформация для удобства
transforms=unwrap
transforms.unwrap.type=io.debezium.transforms.ExtractNewRecordState
# Отправляем только новое состояние записи после изменения
transforms.unwrap.drop.tombstones=true
Эта конфигурация отправляет в Kafka компактные события с итоговым состоянием строки, что упрощает их дальнейшую обработку.