Ответ
Да, я активно использовал CDC для построения пайплайнов данных в реальном времени. На проекте по созданию операционного дашборда нам требовалось отражать изменения из OLTP-базы (PostgreSQL) в аналитическом хранилище с задержкой в секунды.
Реализация:
- Источник: PostgreSQL с включенным логическим декодированием (
wal_level = logical). - CDC-инструмент: Debezium для захвата изменений из WAL-лога PostgreSQL и потоковой передачи их в Apache Kafka. Каждая таблица-источник отображалась на отдельный Kafka-топик.
- Обработка: Потоковый обработчик на Apache Spark Structured Streaming читал топики Kafka. Для каждого события (INSERT, UPDATE, DELETE) применялась бизнес-логика, и результат записывался в целевую витрину данных в S3 (в формате Delta Lake).
Ключевой код (упрощенно):
val df = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "broker:9092")
.option("subscribe", "pg-server.public.orders")
.load()
// Десериализация JSON из Debezium
val ordersChangeData = df.select(from_json(col("value").cast("string"), schema).alias("data"))
.select("data.payload.op", "data.payload.after", "data.payload.before")
// Применение логики (например, upsert в Delta-таблицу)
ordersChangeData.writeStream
.foreachBatch { (batchDF: DataFrame, batchId: Long) =>
batchDF.persist()
// Обработка вставок и обновлений
handleUpserts(batchDF, "op in ('c','u')")
// Обработка удалений
handleDeletes(batchDF, "op = 'd'")
batchDF.unpersist()
}
.start()
Результат: Мы получили актуальную аналитическую картину с минимальной задержкой, что позволило бизнесу быстро реагировать на изменения.