(+351) 21 24 10006  ·  info@bconcepts.pt
Carnaxide, Lisboa

Cómo crear tablas Delta CDC en Lakehouse: paso a paso

João Barros 07 de September de 2026 6 min de lectura

Vamos a implementar un flujo simple de Change Data Capture (CDC) usando tablas Delta en Lakehouse para capturar inserts, updates y deletes de una fuente transaccional. Esto es útil para mantener una sincronización eficiente entre sistemas y permitir análisis históricos sin replicar toda la base de datos. Con un flujo CDC bien diseñado puedes reducir el volumen de datos movidos: por ejemplo, en lugar de reprocesar 100.000 registros por hora, solo aplicas los 1.200 cambios ocurridos en ese periodo.

Requisitos previos

  • Cuenta con acceso a Microsoft Fabric y un Lakehouse configurado.
  • Permisos para crear tablas Delta y ejecutar notebooks con PySpark.
  • Fuente de datos (CSV o stream) con un identificador único y un campo de operación (op: I/U/D) o un log de cambios.

Adicionalmente, es recomendable tener estimaciones de volumen: por ejemplo, saber que esperas en promedio 500–5.000 eventos CDC por hora ayuda a definir micro-batches y la configuración del cluster. Garantizar que tienes políticas de retención y limpieza (VACUUM) alineadas con las necesidades de auditoría es también importante — por ejemplo, retención mínima de 7 días para Delta.

Paso 1: Preparar la tabla Delta destino (del tipo base)

Crea una tabla Delta que servirá como materializado con el estado actual de los registros. Vamos a usar un esquema simple: id, valor, updated_at. Esta tabla mantendrá el último estado por id. Para un escenario real, la tabla puede tener miles a millones de filas; en este ejemplo inicial puedes empezar con 10.000 registros.

from pyspark.sql.types import StructSchema, StructField, IntegerType, StringType, TimestampType

schema = StructType([
    StructField("id", IntegerType(), False),
    StructField("valor", StringType(), True),
    StructField("updated_at", TimestampType(), True)
])

# Path do Lakehouse onde guardar a tabela
target_path = "/lakehouse/yourcatalog/yourdb/delta_target"

# Criar DataFrame vazio e gravar como Delta
empty_df = spark.createDataFrame([], schema)
empty_df.write.format("delta").mode("overwrite").save(target_path)

# Criar tabela no catálogo (opcional)
spark.sql(f"CREATE TABLE IF NOT EXISTS yourdb.delta_target USING DELTA LOCATION '{target_path}'")

Consejo: si esperas muchos registros, considera particionar por una columna adecuada (ej.: ano_mes) o usar Z-order para optimizar lecturas por id.

Paso 2: Ingesta inicial (full load)

Carga el estado inicial de la fuente en la tabla Delta. Esto garantiza que la tabla destino comienza con un snapshot consistente. En un entorno de producción, un full load puede tardar minutos u horas (ej.: 1 millón de registros), por lo que haz esto fuera de horas o con throttling.

# Exemplo a partir de um CSV de origem
source_df = spark.read.format("csv").option("header", True).schema(schema).load("/data/initial_snapshot.csv")

# Gravar no target (substitui conteúdo inicial)
source_df.write.format("delta").mode("overwrite").save(target_path)

Tras el full load valida: por ejemplo, confirma que existen 10.000 registros con spark.read.format("delta").load(target_path).count() o ejecuta checksums por muestreo para garantizar integridad.

Paso 3: Estructurar el feed CDC

Se asume un feed con columnas: id, valor, op, event_time. op = 'I'|'U'|'D'. Prepara un DataFrame con los cambios a aplicar. En una ventana de 1 hora puedes tener, por ejemplo, 300 updates, 150 inserts y 50 deletes — estos números ayudan a definir el tamaño de los batches y la memoria necesaria.

cdc_schema = StructType([
    StructField("id", IntegerType(), False),
    StructField("valor", StringType(), True),
    StructField("op", StringType(), False),
    StructField("event_time", TimestampType(), True)
])

cdc_df = spark.read.format("csv").option("header", True).schema(cdc_schema).load("/data/cdc_feed.csv")

Valida el orden y la exactitud de los timestamps. Si el feed llega desordenado, usa una estrategia de deduplicación basada en event_time.

Paso 4: Aplicar CDC con MERGE (upsert/delete)

Usa MERGE INTO en Delta para aplicar inserts, updates y deletes de forma transaccional. Filtra por orden de eventos si es necesario (ej.: último evento por id). El MERGE garantiza atomicidad — es decir, o todos los cambios del batch entran con éxito, o ninguno entra.

from delta.tables import DeltaTable

delta_table = DeltaTable.forPath(spark, target_path)

# Optional: se existirem múltiplos eventos por id, pega o último por event_time
from pyspark.sql import Window
from pyspark.sql.functions import row_number

w = Window.partitionBy("id").orderBy(cdc_df["event_time"].desc())
cdc_latest = cdc_df.withColumn("rn", row_number().over(w)).filter("rn = 1").drop("rn")

# Executar MERGE
(delta_table.alias("t")
  .merge(cdc_latest.alias("s"), "t.id = s.id")
  .whenMatchedUpdate(
      condition = "s.op != 'D'",
      set = {"valor": "s.valor", "updated_at": "s.event_time"}
  )
  .whenMatchedDelete(condition = "s.op = 'D'")
  .whenNotMatchedInsert(values = {"id": "s.id", "valor": "s.valor", "updated_at": "s.event_time"})
  .execute())

Ejemplo práctico: en un batch con 500 eventos deduplicados, esperas aplicar ~350 updates, 120 inserts y 30 deletes. Tras el MERGE, valida el número de filas y algunos registros aleatorios para confirmar el comportamiento.

Paso 5: Automatizar y tratar errores comunes

Pon este flujo en un notebook programado o en una pipeline. Maneja errores comunes: datos duplicados en el feed, falta del campo op, conflictos de timestamp. Implementa logs y métricas (n.º aplicados, n.º rechazados) para monitorización. Si usas Microsoft Fabric, programa el notebook con triggers y captura fallos para re-procesamiento.

# Exemplo simples de validação antes do MERGE
invalid_ops = cdc_df.filter("op NOT IN ('I','U','D')").count()
if invalid_ops > 0:
    raise ValueError(f"Encontradas {invalid_ops} operações inválidas no feed CDC")

Otras prácticas: usar checkpoints para structured streaming, limitar el tamaño de los batches a 100k eventos, y configurar VACUUM con retención mínima (ej.: 7 días) para evitar la eliminación prematura de historiales necesarios para replays.

Verificar el resultado

Confirma que los registros son correctos consultando la tabla Delta y comparando con el feed esperado. Verifica inserts, updates y deletes y el campo updated_at. Ejemplos de verificación: contar diferencias entre el snapshot esperado y la tabla (por id), comprobar timestamps más recientes y confirmar que no existen filas marcadas como eliminadas cuando no deberían.

# Consultar a tabela Delta
spark.read.format("delta").load(target_path).show()

# Comparar contagens
print("Total na tabela:", spark.read.format("delta").load(target_path).count())

Conclusión

Has implementado un flujo básico de CDC con tablas Delta en Lakehouse, usando MERGE para garantizar atomicidad. Próximos pasos: integrar un stream (structured streaming) para reducir latencia, almacenar metadatos de lineage e implementar pruebas automáticas y alertas. Consejo práctico: comienza con lotes pequeños de CDC (ej.: 1k–5k eventos) y valida siempre el orden de los eventos antes de escalar a producción.