Cómo crear una Delta Live Table en Lakehouse: paso a paso
Este tutorial muestra cómo crear una Delta Live Table en Lakehouse para ingesta continua y transformación de datos, útil cuando necesitas pipelines fiables con historial y monitorización. La tarea se centra en un ejemplo concreto: leer archivos JSON de un directorio, aplicar transformaciones simples y publicar una tabla Delta lista para consumo.
Requisitos previos
- Cuenta con acceso a Microsoft Fabric y permisos para crear un Lakehouse y pipelines.
- Un Lakehouse con un directorio para archivos de ingesta (ej.: /lakehouse/raw/events/).
- Noções básicas de Python/PySpark y del formato Delta.
- Ejemplo de archivos JSON con esquema consistente.
Paso 1: Entender qué es una Delta Live Table en Lakehouse
Delta Live Table es un enfoque para construir pipelines declarativos que producen tablas Delta gestionadas. Permite ingesta continua, manejo de esquemas y monitorización integrada. Este paso es conceptual: visualiza el flujo — origen (archivos JSON) → transformación (limpieza, tipos) → destino (tabla Delta).
Paso 2: Crear un workspace de pipeline
Crea un pipeline en el entorno de Fabric que ejecute código PySpark en trigger continuo o programado. Define el nombre, el tipo de runtime y el Lakehouse donde se grabarán las tablas.
# Exemplo conceptual (interface do Fabric tem gui); se usares CLI/SDK configura o pipeline com:
# runtime: pyspark
# target: Lakehouse//tables/
Paso 3: Código PySpark mínimo para una Delta Live Table
Escribe un script PySpark que lea JSON en modo streaming, aplique transformaciones simples y escriba en una tabla Delta. Mantén el código mínimo para probar el flujo.
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, to_timestamp
spark = SparkSession.builder.getOrCreate()
# Fonte: ficheiros JSON em modo 'cloudFiles' (auto-infer schema quando disponível)
raw_df = (
spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.inferColumnTypes", "true")
.load("/lakehouse/raw/events/")
)
# Transformações: limpar campos e converter timestamps
clean_df = (
raw_df
.withColumn("event_time", to_timestamp(col("event_time"), "yyyy-MM-dd'T'HH:mm:ss"))
.withColumn("user_id", col("user_id").cast("string"))
.na.drop(subset=["event_time", "user_id"]) # descartar linhas inválidas
)
# Escrita em Delta para o Lakehouse (modo append, checkpoint obrigatório)
(output = clean_df.writeStream
.format("delta")
.outputMode("append")
.option("checkpointLocation", "/lakehouse/checkpoints/events_checkpoint/")
.start("/lakehouse/tables/events_delta/"))
Paso 4: Configurar políticas de esquema y manejo de errores
Define cómo el pipeline gestiona los cambios de esquema y los datos corruptos. En entornos Fabric/Lakehouse, configura opciones de Auto Loader o cloudFiles para ignorar archivos corruptos y registrar errores.
# Opções adicionais para robustez (exemplo cloudFiles)
.option("cloudFiles.schemaEvolutionMode", "addNewColumns")
.option("cloudFiles.maxFilesPerTrigger", "100")
.option("badRecordsPath", "/lakehouse/errors/bad_records/")
Paso 5: Validar y crear la tabla Delta gestionada
Después de que el stream se escriba en la ruta Delta, registra la carpeta como tabla Delta en el catálogo del Lakehouse para consulta SQL y consumo por Power BI.
-- SQL ejecutado no SQL endpoint ou notebook para crear tabela a partir do caminho Delta
CREATE TABLE IF NOT EXISTS lakehouse.events_delta
USING DELTA
LOCATION '/lakehouse/tables/events_delta/';
Paso 6: Operaciones comunes y errores a evitar
Monitorea checkpointLocation para evitar pérdida de estado; no compartas el mismo checkpoint entre pipelines distintos. Si encuentras error de schema mismatch, activa schemaEvolution o realiza validación previa. Evita escribir directamente en el mismo directorio con jobs concurrentes sin ACID (Delta resuelve parcialmente, pero tiene límites).
# Erros comuns (consulta de diagnóstico)
-- Verifica o estado da tabela
DESCRIBE HISTORY lakehouse.events_delta;
-- Conferir ficheiros corruptos
ls /lakehouse/errors/bad_records/
Verificar el resultado
Confirma que la tabla Delta recibió datos: consulta la tabla con SQL y verifica que hay datos recientes. Comprueba la existencia del checkpoint y los archivos _delta_log en el directorio de la tabla.
SELECT COUNT(*) FROM lakehouse.events_delta;
SELECT event_time, user_id FROM lakehouse.events_delta ORDER BY event_time DESC LIMIT 10;
# No storage, confirma:
# /lakehouse/tables/events_delta/_delta_log/
# /lakehouse/checkpoints/events_checkpoint/
Conclusión
Ahora tienes un pipeline básico de Delta Live Table en Lakehouse para ingesta continua de JSON, transformación mínima y publicación como tabla Delta. Próximos pasos: añadir pruebas de calidad, enriquecimientos y monitorización en Fabric. Consejo: comienza con un conjunto pequeño de archivos para validar el esquema antes de activar el streaming en producción.