Работали ли вы с CDC (Change Data Capture)?

«Работали ли вы с CDC (Change Data Capture)?» — вопрос из категории ETL и пайплайны данных, который задают на 33% собеседований Data Инженер. Ниже — развёрнутый ответ с разбором ключевых моментов.

Ответ

Да, я активно использовал CDC для построения пайплайнов данных в реальном времени. На проекте по созданию операционного дашборда нам требовалось отражать изменения из OLTP-базы (PostgreSQL) в аналитическом хранилище с задержкой в секунды.

Реализация:

  1. Источник: PostgreSQL с включенным логическим декодированием (wal_level = logical).
  2. CDC-инструмент: Debezium для захвата изменений из WAL-лога PostgreSQL и потоковой передачи их в Apache Kafka. Каждая таблица-источник отображалась на отдельный Kafka-топик.
  3. Обработка: Потоковый обработчик на 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()

Результат: Мы получили актуальную аналитическую картину с минимальной задержкой, что позволило бизнесу быстро реагировать на изменения.