Trabajar con tablas Delta Lake en Microsoft Fabric
El módulo más importante de la Fase 1: la tecnología que sostiene todo el data engineering en Fabric. MERGE, SCD, schema evolution, time travel, OPTIMIZE, V-Order, VACUUM y streaming. 🌸
IntermedioAvanzado
🎯
1. Objetivos y encaje en el DP-700
¿Qué te enseña este módulo?
Este es el módulo más importante de toda la Fase 1. Aquí aprendemos a fondo la tecnología que sostiene absolutamente todo el data engineering en Fabric. Los objetivos oficiales:
Entender Delta Lake y las Delta tables en Microsoft Fabric.
Crear y gestionar tablas Delta usando Spark.
Optimizar tablas Delta.
Usar Spark para consultar y transformar datos en Delta tables.
Usar Delta tables con Spark Structured Streaming.
Peso en el examen DP-700
Delta Lake aparece transversalmente en todo el examen, pero especialmente en:
Dominio
Cómo aparece Delta
Implement and manage (30-35%)
Diseño de tablas Delta en arquitecturas medallion, permisos sobre tablas
Ingest and transform (30-35%)
MERGE, upsert, SCD, deduplicación, late-arriving data — todo con Delta
Monitor and optimize (30-35%)
OPTIMIZE, VACUUM, V-Order, particionado, streaming Delta — todo el mantenimiento
🎯 Qué es crítico para dominar:
MERGE INTO para upserts, SCD Type 1 y Type 2
Time travel (versionAsOf, timestampAsOf, DESCRIBE HISTORY)
OPTIMIZE + V-Order + VACUUM para performance y limpieza
Streaming con Delta como source/sink (Structured Streaming)
Schema evolution para gestionar cambios de estructura
💎
2. ¿Qué es Delta Lake exactamente?
Definición
Delta Lake es una storage layer open-source que añade semántica de base de datos relacional al procesamiento data lake basado en Spark.
En cristiano: es el formato que hace que un montón de ficheros Parquet se comporten como una tabla de base de datos con transacciones, schema, versiones, etc. 🌸
Dónde ves que una tabla es Delta
En el Lakehouse explorer de Fabric, las tablas Delta se identifican por el icono triangular delta (Δ) al lado del nombre. Si no ves ese icono, esa "tabla" no es Delta y no tendrá las capacidades avanzadas.
Delta en el ecosistema Fabric
En Fabric, todas las tablas de un lakehouse son Delta tables por defecto. No tienes que hacer nada especial: cuando escribes datos a Tables/, se guardan en formato Delta automáticamente.
🌷 Analogía kawaii: piensa en Parquet como hojas sueltas de un libro. Delta Lake es como encuadernar esas hojas con un índice detallado y un registro de cambios. Ahora ese conjunto de hojas es un libro que puedes actualizar, versionar y auditar. 📖✨
🩺
3. Anatomía de una tabla Delta
Estructura física en OneLake
Cada tabla Delta ocupa una carpeta con esta estructura:
Tables/mi_tabla_delta/
├── part-00000-abc123....snappy.parquet ← Datos (Parquet)
├── part-00001-def456....snappy.parquet
├── part-00002-ghi789....snappy.parquet
├── ...
└── _delta_log/ ← Transaction log
├── 00000000000000000000.json ← Cada operación = 1 JSON
├── 00000000000000000001.json
├── 00000000000000000002.json
├── ...
└── 00000000000000000010.checkpoint.parquet ← Checkpoint cada N versiones
Los dos ingredientes
A) Ficheros Parquet (datos)
Formato columnar comprimido.
Inmutables: nunca se modifican en sitio. Cada cambio genera nuevos ficheros.
Contienen los valores reales de las filas.
B) Carpeta _delta_log/ (metadatos)
JSON por cada transacción: describe qué ficheros forman parte de la tabla en cada versión.
Contiene información de add, remove, schema, metadata.
Cada N versiones se genera un checkpoint en Parquet (para acelerar la lectura del log).
🎯 Cómo Delta sabe qué ficheros están "activos"
Cuando lees una tabla Delta:
Delta lee el último estado consolidado del _delta_log.
Ese estado dice: "los ficheros activos son X, Y, Z".
Delta lee solo esos ficheros.
Los ficheros viejos siguen físicamente ahí (hasta que hagas VACUUM) → habilitan time travel.
¿Cómo se hace UPDATE si Parquet es inmutable?
Delta no modifica en sitio. Cuando haces UPDATE:
Delta identifica los ficheros que contienen las filas afectadas.
Lee esos ficheros.
Aplica los cambios en memoria.
Escribe nuevos ficheros con las filas actualizadas.
Marca los ficheros viejos como removidos en el nuevo JSON del _delta_log.
Nuevas lecturas ven la nueva versión.
Es "copy-on-write". Por eso a veces se acumulan muchos ficheros pequeños → problema del small files, que resolvemos con OPTIMIZE.
✨
4. Los 5 grandes beneficios de Delta
Estos son los 5 puntos oficiales de la doc de Microsoft. Apréndetelos porque salen en preguntas del tipo "¿Cuáles son los beneficios de usar Delta?".
🎯 A) Tablas relacionales con operaciones CRUD
Con Spark puedes:
CREATE tablas Delta
SELECT filas
INSERT nuevas filas
UPDATE filas existentes
DELETE filas
MERGE (upsert)
Es decir, como una BBDD relacional, pero encima de tu data lake.
🎯 B) Transacciones ACID
Delta implementa las 4 propiedades ACID sobre Spark:
Atomicity: una transacción se completa entera o no se aplica nada.
Consistency: la tabla siempre queda en un estado válido.
Isolation: procesos concurrentes no interfieren entre sí (serializable isolation).
Durability: los cambios persisten aunque falle el sistema.
Esto es lo que permite que varios notebooks escriban a la misma tabla simultáneamente sin corromperla.
🎯 C) Versionado de datos y time travel
Como todas las transacciones se registran, puedes:
Ver el historial completo de cambios.
Consultar la tabla como estaba en un momento pasado.
Auditar cambios.
Recuperarte de errores accidentales.
Ejemplos: versionAsOf, timestampAsOf, DESCRIBE HISTORY.
🎯 D) Soporte para batch Y streaming
Delta puede ser:
Sink (destino) de un stream → los datos que llegan se acumulan en la tabla.
Source (fuente) de un stream → cambios en la tabla se procesan como stream.
Con la misma API que para batch. Esto es súper potente para escenarios de IoT, logs, telemetría.
🎯 E) Formato estándar e interoperabilidad
Delta usa Parquet por debajo, que es el formato estándar en data lakes. Esto significa:
Cualquier motor Parquet puede leerlo.
El SQL analytics endpoint puede consultar tablas Delta con T-SQL.
Interoperabilidad con Databricks, Synapse, Snowflake y otros que soporten Delta.
🛠️
5. Crear tablas Delta: 5 formas distintas
Hay múltiples caminos para crear una tabla Delta. Cada uno tiene su caso de uso.
Método 1: saveAsTable() desde un DataFrame (el más común)
# Cargar datos en un DataFrame
df = spark.read.load('Files/mydata.csv', format='csv', header=True)
# Guardarlos como tabla Delta gestionada
df.write.format("delta").saveAsTable("mytable")
Qué pasa:
Los datos se guardan como Parquet en Tables/mytable/
Se crea el _delta_log/
La tabla aparece en el Lakehouse Explorer bajo Tables
Es una tabla managed (Fabric gestiona todo)
Método 2: saveAsTable() con path para external table
# Tabla external: metadata en el metastore, datos en Files/
df.write.format("delta").saveAsTable(
"myexternaltable",
path="Files/myexternaltable"
)
Uso típico: cuando quieres preparar el schema primero y luego llenar la tabla con datos que vienen por otro camino (streaming, notebooks múltiples, etc.).
Método 4: SQL CREATE TABLE
%%sql
CREATE TABLE salesorders
(
OrderId INT NOT NULL,
OrderDate TIMESTAMP NOT NULL,
CustomerName STRING,
SalesTotal FLOAT NOT NULL
)
USING DELTA
Sintaxis SQL clásica. Para external:
%%sql
CREATE TABLE MyExternalTable
USING DELTA
LOCATION 'Files/mydata'
Cuando creas una external table con CREATE TABLE ... LOCATION, el schema se infiere de los ficheros Parquet ya presentes en esa ubicación.
Método 5: Escribir Delta sin crear tabla en metastore
A veces quieres guardar datos en formato Delta pero sin registrar la tabla en el metastore:
Escribe ficheros Delta a esa ruta (con su _delta_log/), pero NO aparece como tabla en el catálogo.
⚠️ Tip curioso: si usas esta técnica y guardas en la ruta Tables/ en vez de Files/, Fabric usa automatic table discovery y crea el metadata automáticamente. Bam, tabla registrada sin saveAsTable. 🌸
Modos de escritura (.mode())
Aplican a save() y saveAsTable():
# Sobrescribir cualquier fichero existente
df.write.format("delta").mode("overwrite").save(delta_path)
# Añadir filas nuevas a lo existente
df.write.format("delta").mode("append").save(delta_path)
Modo
Descripción
overwrite
Borra lo existente y escribe de nuevo
append
Añade a lo existente
ignore
Si existe, no hace nada
error (default)
Falla si ya existe
Comparativa rápida de métodos
Método
Cuándo usarlo
df.write.saveAsTable("t")
Tabla managed desde datos existentes
df.write.saveAsTable("t", path=...)
External table desde datos existentes
DeltaTable.create()
Tabla vacía con schema definido
CREATE TABLE ... USING DELTA
SQL clásico (managed o external)
df.write.save(path)
Datos Delta sin metadata en catálogo
🔧
6. Managed vs External (repaso profundo)
Este tema ya lo vimos en módulos anteriores, pero aquí lo vemos desde la perspectiva Delta.
Managed table
Definición: Fabric controla metadata + datos.
Ubicación datos: Tables/<nombre_tabla>/ en el lakehouse.
Creación: sin especificar path.
DROP: borra metadata Y datos.
# Managed
df.write.format("delta").saveAsTable("managed_products")
# Al dropear...
spark.sql("DROP TABLE managed_products")
# ✅ Se borra el metadata
# ✅ Se borran los ficheros Parquet + _delta_log de Tables/managed_products/
External table
Definición: Fabric controla solo el metadata; los datos viven donde tú digas.
Ubicación datos: la ruta que especificas con path o LOCATION.
Creación: con path=... o LOCATION '...'.
DROP: borra solo el metadata. Los datos permanecen.
# External
df.write.format("delta").saveAsTable(
"external_products",
path="Files/external_products"
)
# Al dropear...
spark.sql("DROP TABLE external_products")
# ✅ Se borra el metadata
# ❌ Los ficheros Parquet + _delta_log de Files/external_products/ permanecen
🎯 Cuando el schema se infiere
Cuando creas una external table con LOCATIONapuntando a datos ya existentes:
%%sql
CREATE TABLE MyExternalTable
USING DELTA
LOCATION 'Files/mydata'
El schema se infiere de los ficheros Parquet en Files/mydata/. No hace falta especificarlo.
⚠️ Recomendación oficial en Fabric. Recordamos del Módulo 2:
"Lakehouses support only managed tables. External tables with CREATE TABLE ... LOCATION aren't supported. To access Delta tables in other storage locations, use shortcuts instead."
Managed tables → totalmente soportadas y recomendadas ✅
External tables con SQL DDL puro (CREATE TABLE ... LOCATION) → no oficialmente soportadas, pueden dar problemas
External tables con df.write.saveAsTable(name, path=...) → funcionan técnicamente
Para acceder a datos externos → shortcuts ← 🎯 esta es LA respuesta esperada en el examen
🔨
7. DeltaTableBuilder API
Esta API te permite crear tablas Delta programáticamente con control fino.
# Solo crea si no existe (no falla si ya está)
DeltaTable.createIfNotExists(spark) \
.tableName("products") \
.addColumn("id", "INT") \
.execute()
# Sobrescribe si existe
DeltaTable.createOrReplace(spark) \
.tableName("products") \
.addColumn("id", "INT") \
.execute()
Cuándo usarlo vs SQL
DeltaTableBuilder API: cuando construyes el schema programáticamente (por ejemplo, generado dinámicamente), o cuando tienes lógica condicional.
SQL CREATE TABLE: cuando conoces el schema de antemano y es más legible en SQL.
Ambos son equivalentes en cuanto a lo que producen.
💾
8. Guardar datos en formato Delta sin tabla
A veces quieres persistir datos en formato Delta sin registrarlos como tabla. Usos típicos:
Datos intermedios de un ETL que no necesitas exponer.
Datos que consumirás con la Delta API directamente por ruta.
Datos que luego "overlayarás" con una tabla definida.
# Sin necesidad de que exista en el metastore
df = spark.read.format("delta").load("Files/mydatatable")
display(df)
Usar la Delta API por ruta
from delta.tables import *
# Referenciar tabla Delta por ruta
deltaTable = DeltaTable.forPath(spark, "Files/mydatatable")
# Operar sobre ella
deltaTable.update(
condition="Category = 'Discontinued'",
set={"Active": "false"}
)
🎯 Tip: Automatic table discovery
Si guardas datos Delta en Tables/nombre/, Fabric detecta automáticamente y crea el metadata:
# Los datos van a Files (no se registra como tabla)
df.write.format("delta").save("Files/mydatatable") # No aparece en Tables
# Los datos van a Tables (Fabric detecta y registra automáticamente)
df.write.format("delta").save("Tables/mydatatable") # ✨ Aparece como tabla
🔄
9. Operaciones de datos: SELECT, INSERT, UPDATE, DELETE
Aquí es donde Delta brilla: puedes hacer operaciones relacionales completas sobre tu data lake.
SELECT (leer)
Spark SQL:
%%sql
SELECT * FROM products WHERE Category = 'Bikes'
PySpark con spark.sql:
df = spark.sql("SELECT * FROM products WHERE Category = 'Bikes'")
display(df)
MERGE INTO es probablemente la operación más importante de Delta Lake para casos reales. Combina INSERT, UPDATE y DELETE en una sola operación atómica basada en si las filas hacen match o no.
Sintaxis básica
MERGE INTO target_table t
USING source_table s
ON t.key = s.key
WHEN MATCHED THEN UPDATE SET t.value = s.value
WHEN NOT MATCHED THEN INSERT (key, value) VALUES (s.key, s.value)
En cristiano:
MERGE INTO target USING source → "Voy a fusionar filas de source en target"
ON condition → "Considera que hacen match si..."
WHEN MATCHED → "Cuando una fila hace match..."
WHEN NOT MATCHED → "Cuando una fila del source NO existe en target..."
WHEN NOT MATCHED BY SOURCE → "Cuando una fila del target NO está en source..." (útil para deletes)
Caso de uso 1: Upsert clásico
Tienes una tabla people10m (target). Recibes actualizaciones diarias en people10m_updates (source). Quieres actualizar los que existen e insertar los nuevos.
MERGE INTO people10m t
USING people10m_updates s
ON t.id = s.id
WHEN MATCHED THEN UPDATE SET * -- Actualiza todas las columnas
WHEN NOT MATCHED THEN INSERT * -- Inserta todas las columnas
El * funciona cuando ambos tienen las mismas columnas.
Caso de uso 2: MERGE con lógica condicional
MERGE INTO products t
USING product_updates s
ON t.ProductId = s.ProductId
WHEN MATCHED AND s.Category = 'Discontinued' THEN
DELETE
WHEN MATCHED AND s.Price != t.Price THEN
UPDATE SET t.Price = s.Price, t.UpdatedAt = current_timestamp()
WHEN NOT MATCHED AND s.Category IS NOT NULL THEN
INSERT (ProductId, Name, Category, Price)
VALUES (s.ProductId, s.Name, s.Category, s.Price)
Aquí demostramos:
MERGE que hace DELETE si la fila está discontinued.
MERGE que hace UPDATE solo si el precio cambió.
MERGE que hace INSERT solo si tiene categoría.
Caso de uso 3: Deduplicación
Insertar filas nuevas de un source, ignorando las duplicadas:
MERGE INTO orders t
USING new_orders s
ON t.OrderId = s.OrderId
WHEN NOT MATCHED THEN
INSERT *
Como no hay WHEN MATCHED, si el OrderId ya existe → no hace nada. Los duplicados se descartan silenciosamente. Muy útil para exactly-once semantics en pipelines de ingesta.
Caso de uso 4: WHEN NOT MATCHED BY SOURCE (deletes)
Quieres que target refleje exactamente lo que hay en source (útil para sincronización total):
MERGE INTO customers t
USING (SELECT * FROM customers_source WHERE updated_at >= current_date() - INTERVAL '5' DAY) s
ON t.id = s.id
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *
WHEN NOT MATCHED BY SOURCE
AND t.updated_at >= current_date() - INTERVAL '5' DAY
THEN DELETE
Este pattern es genial para escenarios donde el source puede cambiar los últimos N días y tú quieres reflejar exactamente esos cambios en target.
🎯 Caso de uso 5: SCD Type 1 (sobrescribir historial)
SCD Type 1: cada cambio sobrescribe el valor anterior, sin guardar historia. La tabla siempre refleja el estado actual.
MERGE INTO dim_customer t
USING staging_customer s
ON t.CustomerKey = s.CustomerKey
WHEN MATCHED THEN
UPDATE SET
t.Name = s.Name,
t.Email = s.Email,
t.State = s.State,
t.UpdatedAt = current_timestamp()
WHEN NOT MATCHED THEN
INSERT (CustomerKey, Name, Email, State, CreatedAt)
VALUES (s.CustomerKey, s.Name, s.Email, s.State, current_timestamp())
Fácil, eficiente. Cuando un cliente cambia de state, se sobrescribe.
🎯 Caso de uso 6: SCD Type 2 (guardar historial)
SCD Type 2: cada cambio crea una nueva fila con validez temporal (Valid_From, Valid_To, Is_Current). La tabla guarda TODAS las versiones históricas.
Ejemplo del comportamiento:
CustomerKey
CustomerID
Name
State
Valid_From
Valid_To
Is_Current
1001
C-123
Company
CA
2023-01-15
2026-02-20
No
1002
C-123
Company
NY
2026-02-20
NULL
Yes
Implementación con MERGE en dos pasos.
Paso 1: cerrar las filas actuales que han cambiado.
MERGE INTO dim_customer_scd2 t
USING (
SELECT s.CustomerID,
s.Name,
s.State,
current_date() AS effective_date
FROM staging_customer s
) s
ON t.CustomerID = s.CustomerID AND t.Is_Current = true AND t.State != s.State
WHEN MATCHED THEN
UPDATE SET
t.Valid_To = s.effective_date,
t.Is_Current = false
Paso 2: insertar las nuevas versiones.
INSERT INTO dim_customer_scd2 (CustomerID, Name, State, Valid_From, Valid_To, Is_Current)
SELECT
s.CustomerID, s.Name, s.State,
current_date(), NULL, true
FROM staging_customer s
LEFT JOIN dim_customer_scd2 t
ON s.CustomerID = t.CustomerID AND t.Is_Current = true
WHERE t.CustomerID IS NULL -- No existía
OR t.State != s.State -- O ha cambiado
⚠️ Ojo: SCD Type 2 con MERGE puro es tricky por las semánticas de match. En Fabric, Copy Job soporta SCD Type 2 nativamente (preview), simplificando esto.
🎯 Limitación importante
"A merge operation can fail if multiple rows of the source dataset match and the merge attempts to update the same rows of the target Delta Lake table."
Es decir: si en el source hay filas duplicadas que hacen match con la misma fila del target, MERGE falla con error (por ambigüedad semántica).
Solución: deduplicar el source antes de MERGE:
source_dedup = source_df.dropDuplicates(["id"])
SCD Type 1 vs SCD Type 2 (comparativa oficial)
Feature
SCD Type 1 (Merge)
SCD Type 2
Soporte en Copy Job de Fabric
✅
✅ (Preview)
Historia
No preservada
Preservada con filas versionadas
Estado destino
Refleja siempre el source actual
Contiene todas las versiones
Deletes
Filas físicamente eliminadas
Soft delete (Is_Current = false)
Caso de uso
Reporting operacional, sync realtime
Análisis histórico, auditoría, compliance
Complejidad
Baja
Alta
🔀
12. Schema evolution
Schema evolution = capacidad de modificar el schema de una tabla Delta sin romper compatibilidad.
Por qué importa
En la vida real, los schemas cambian:
Añades columnas nuevas.
Cambias tipos de datos.
Renombras (menos común).
Delta permite gestionar estos cambios con seguridad.
Añadir columnas al escribir
Con mergeSchema en escritura:
# El source tiene columnas nuevas que no están en target
new_df.write.format("delta") \
.mode("append") \
.option("mergeSchema", "true") \
.saveAsTable("products")
Las columnas nuevas del source se añaden automáticamente al target. Las filas existentes tendrán NULL en esas columnas.
Sin mergeSchema → falla con error de schema mismatch (schema enforcement).
⚠️ Cuidado: esto sobrescribe todo el schema. Úsalo solo cuando estés segura.
Schema evolution con MERGE
Con MERGE, puedes usar WITH SCHEMA EVOLUTION:
MERGE WITH SCHEMA EVOLUTION INTO sales AS target
USING staged_sales AS source
ON target.order_id = source.order_id
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *
Si staged_sales tiene columnas nuevas, se añaden automáticamente a sales.
Cuando esté activo, cualquier MERGE hará schema evolution automáticamente.
ALTER TABLE (cambios explícitos)
%%sql
-- Añadir columna
ALTER TABLE products ADD COLUMNS (Discount FLOAT)
-- Cambiar comment
ALTER TABLE products ALTER COLUMN Name COMMENT 'Product name'
-- Renombrar (requiere propiedad especial habilitada)
ALTER TABLE products RENAME COLUMN OldName TO NewName
Limitación: schema enforcement por defecto
Por defecto, Delta rechaza escrituras que no cumplen el schema:
Columna de más en source → error
Tipo distinto → error
Columna faltante → puede o no fallar según el modo
📌 Buena práctica: el schema enforcement por defecto es lo correcto. Activa schema evolution solo cuando realmente quieras cambios.
⏰
13. Time travel y DESCRIBE HISTORY
Uno de los superpoderes más chulos de Delta. 🌸
Ver el historial de una tabla
Por nombre:
%%sql
DESCRIBE HISTORY products
Por ruta (external tables):
%%sql
DESCRIBE HISTORY 'Files/mytable'
Resultado típico:
version
timestamp
operation
operationParameters
2
2026-04-04T21:46:43Z
UPDATE
{"predicate":"(ProductId = 1)"}
1
2026-04-04T21:42:48Z
WRITE
{"mode":"Append","partitionBy":"[]"}
0
2026-04-04T20:04:23Z
CREATE TABLE
{"isManaged":"true"}
Ves cada versión, el timestamp, la operación y los parámetros.
Consultar una versión específica (por número)
# Leer la tabla como estaba en la versión 0
df = spark.read.format("delta") \
.option("versionAsOf", 0) \
.load("Files/mytable")
display(df)
O con SQL:
%%sql
SELECT * FROM products VERSION AS OF 0
Consultar por timestamp
# La tabla como estaba en enero 2026
df = spark.read.format("delta") \
.option("timestampAsOf", "2026-01-01") \
.load("Files/mytable")
O con SQL:
%%sql
SELECT * FROM products TIMESTAMP AS OF '2026-01-01T00:00:00'
Casos de uso reales
🔍 Auditoría: "¿cómo estaba esta tabla el día antes del incidente?"
🔄 Revertir un error: "reescribe la tabla desde la versión X"
📊 Análisis histórico: comparar estado actual con estado hace un mes.
🐛 Debugging: reproducir un bug que se dio hace días.
🎯 Rollback: revertir a una versión anterior
%%sql
-- Restaurar la tabla al estado de la versión 3
RESTORE TABLE products TO VERSION AS OF 3
O:
%%sql
RESTORE TABLE products TO TIMESTAMP AS OF '2026-04-01'
⚠️ Requiere que los ficheros de esa versión aún existan (no hayan sido purgados por VACUUM).
⚠️ Limitación crítica: retención. Time travel funciona mientras existan los ficheros Parquet viejos. Los ficheros viejos se purgan con VACUUM después del periodo de retención (default 7 días). Después de VACUUM, ya no puedes hacer time travel a versiones más viejas que el threshold.
Puedes cambiar la retención:
ALTER TABLE products SET TBLPROPERTIES ('delta.logRetentionDuration' = 'interval 30 days')
🎯 Propiedades relacionadas
Propiedad
Descripción
Default
delta.logRetentionDuration
Cuánto se guarda el transaction log
30 días
delta.deletedFileRetentionDuration
Cuánto se guardan los ficheros marcados como removidos
7 días
Después de estos periodos, VACUUM puede limpiar.
⚡
14. Optimización: OptimizeWrite, OPTIMIZE y V-Order
Tres features clave para mantener el rendimiento. Este apartado es súper examinable.
El problema: "small files problem"
Cada operación de escritura en Delta genera nuevos ficheros Parquet. Si haces muchas escrituras pequeñas (por ejemplo, streaming a alta frecuencia) → acumulas muchísimos ficheros pequeños.
Consecuencia:
Lecturas más lentas (Spark tiene que abrir muchos ficheros).
Queries pueden fallar por overhead.
Costes más altos.
Delta tiene 3 features para resolverlo.
🎯 OptimizeWrite (prevención)
OptimizeWrite es una feature que al escribir, consolida datos en menos ficheros más grandes. Previene el problema antes de que ocurra.
En Fabric está activada por defecto.
# Desactivar
spark.conf.set("spark.microsoft.delta.optimizeWrite.enabled", False)
# Activar
spark.conf.set("spark.microsoft.delta.optimizeWrite.enabled", True)
# Consultar estado actual
print(spark.conf.get("spark.microsoft.delta.optimizeWrite.enabled"))
OPTIMIZE es un comando de mantenimiento que consolida ficheros Parquet pequeños en ficheros más grandes.
Cuándo ejecutar:
Después de cargar tablas grandes.
Periódicamente en tablas con muchas escrituras.
Cuando notas queries lentas.
Beneficios: menos ficheros, mejor compresión, mejor distribución en nodos y queries más rápidas.
Cómo ejecutar (3 formas):
A) Desde el Lakehouse Explorer (UI): menú ... junto a la tabla → Maintenance → Run OPTIMIZE command → opcionalmente marcar V-Order → Run now.
B) Desde SQL:
%%sql
OPTIMIZE products
C) Con Z-ORDER (co-localizar datos por columna):
%%sql
OPTIMIZE products ZORDER BY (Category)
ZORDER es una técnica avanzada que reordena los datos dentro de los ficheros para acelerar filtros por esa columna.
🎯 V-Order (aceleración de lecturas)
V-Order es una feature exclusiva de Fabric que optimiza los ficheros Parquet para lecturas ultra rápidas desde motores de Fabric.
Características clave:
✅ Activada por defecto en Fabric.
⚡ Habilita "lightning-fast reads" con acceso in-memory-like.
💰 Reduce recursos de red, disco y CPU en lecturas.
📈 Aplicada al escribir datos.
🐢 Overhead pequeño (~15%) en escritura, gran beneficio en lectura.
📊 Compatible con motores externos (100% Parquet compliant).
Cómo funciona técnicamente: aplica sorting especial, distribución de row groups optimizada, dictionary encoding y compresión.
Efecto en distintos motores:
Motor
Beneficio
Power BI (VertiScan)
Máximo aprovechamiento
SQL analytics endpoint (VertiScan)
Máximo aprovechamiento
Spark
~10% más rápido (hasta 50%)
Otros motores Parquet
Aún leen sin problema, beneficio moderado
Cuándo desactivarlo: en write-intensive scenarios como staging bronze donde los datos se leen 1-2 veces solo, el overhead del 15% no se compensa. En esos casos, desactivar V-Order acelera la ingesta.
Aplicar V-Order con OPTIMIZE: desde la UI en Table Maintenance → marcar checkbox "Apply V-order"; o automáticamente al escribir (si está activo por defecto).
Comparativa de las 3 features
Feature
Cuándo actúa
Objetivo
Default en Fabric
OptimizeWrite
Al escribir
Prevenir small files
✅ ON
OPTIMIZE
Manual/scheduled
Consolidar ficheros existentes
Manual
V-Order
Al escribir
Acelerar lecturas futuras
✅ ON
🧹
15. VACUUM: limpieza y sus peligros
¿Qué hace VACUUM?
VACUUMelimina físicamente los ficheros Parquet obsoletos de una tabla Delta:
Ficheros que ya no están referenciados en el transaction log.
Ficheros más antiguos que el retention period.
IMPORTANTE: VACUUM NO borra el transaction log, solo los ficheros de datos. El historial de operaciones sigue en el _delta_log/.
🎯 Impacto en time travel: cuando ejecutas VACUUM, ya no puedes hacer time travel a versiones anteriores al retention period. Ejemplo: retention = 7 días, hoy es 2026-05-01, VACUUM borra ficheros de antes de 2026-04-24 → ya no puedes hacer VERSION AS OF a versiones que dependían de esos ficheros.
Retention period
Default: 7 días (168 horas).
El sistema impide por defecto usar un retention menor a 7 días, porque puede causar problemas:
Rompe time travel.
Puede corromper lecturas concurrentes (readers largos que aún referencian ficheros que se están borrando).
Cómo ejecutar VACUUM
Desde el Lakehouse Explorer: menú ... junto a la tabla → Maintenance → Run VACUUM command using retention threshold → configurar threshold → Run now.
Con SQL:
%%sql
VACUUM products RETAIN 168 HOURS
Con lakehouse.tabla:
%%sql
VACUUM lakehouse2.products RETAIN 168 HOURS
Con retention custom (más de 7 días, más seguro):
%%sql
VACUUM products RETAIN 720 HOURS -- 30 días
⚠️ Cómo bajar el retention (peligroso)
Si REALMENTE necesitas bajar de 7 días (raramente):
Puedes dividir en pocas particiones grandes (no miles).
Las queries filtran frecuentemente por esa columna.
NO particionar cuando:
Los volúmenes son pequeños (< 1GB total).
La columna tiene alta cardinalidad (muchos valores únicos).
Necesitarías múltiples niveles de particionado que no aportan.
📌 Regla del pulgar: ~1 GB por partición es un buen tamaño y no más de unos cientos de particiones en total. Ejemplos buenos: year, country, region, category (baja cardinalidad). Ejemplos malos: user_id, order_id, timestamp a segundo.
Limitación importante
Las particiones son un layout fijo. Cuando decides particionar por year, esa decisión afecta todos los patrones de query. No se adapta dinámicamente.
Alternativas modernas para agilidad:
Liquid Clustering (feature reciente de Delta) — más flexible.
Z-ORDER con OPTIMIZE — reorganiza sin cambiar layout físico.
La columna Categorydesaparece del DataFrame (Spark ya sabe que todos son "Road Bikes"). Si haces road_bikes.columns, no verás Category.
🌊
17. Streaming: Delta como source y sink
Una de las características más potentes de Delta: se integra con Spark Structured Streaming para casos de streaming de datos.
¿Qué es Spark Structured Streaming?
Es una API de Spark para procesar streams de datos en tiempo real. La abstracción clave: un stream es un DataFrame ilimitado que va creciendo.
Con Spark Structured Streaming puedes leer de:
Ports de red.
Sistemas de mensajería: Azure Event Hubs, Kafka.
Ubicaciones de file system (poll files).
Delta tables ← la que nos interesa.
Y escribir a Delta tables ← la otra que nos interesa, y muchas otras.
Delta como sink (destino de un stream)
Caso: capturar telemetría IoT y guardarla en una tabla Delta.
# Leer stream desde una fuente (por ejemplo Event Hubs)
stream_df = spark.readStream.format("eventhubs") \
.options(**event_hub_conf) \
.load()
# Escribir a Delta
query = stream_df.writeStream \
.format("delta") \
.option("checkpointLocation", "Files/checkpoints/iot_stream") \
.start("Tables/iot_readings")
# Verificar
print("Streaming to iot_readings...")
Delta como source (leer una tabla como stream)
Caso: procesar en real time los cambios de una tabla Delta.
# Leer una tabla Delta como stream
stream_df = spark.readStream.format("delta") \
.option("ignoreChanges", "true") \
.table("orders_in")
# Verificar que es streaming
print(stream_df.isStreaming) # True
⚠️ Importante: cuando usas una tabla Delta como source, solo se pueden incluir operaciones append en el stream. Modificaciones causan error a menos que uses ignoreChanges = true (ignora cambios que no sean appends) o ignoreDeletes = true (ignora eliminaciones).
Para hacer MERGE (upsert) desde un stream, usas foreachBatch:
def merge_batch(microbatch_df, batch_id):
# microbatch_df es un DataFrame estándar (no streaming)
# Puedes hacer MERGE aquí
microbatch_df.createOrReplaceTempView("updates")
spark.sql("""
MERGE INTO target t
USING updates s
ON t.id = s.id
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *
""")
# Aplicar la función a cada microbatch del stream
stream_df.writeStream \
.foreachBatch(merge_batch) \
.option("checkpointLocation", "Files/checkpoints/merge_stream") \
.start()
⚠️ Idempotencia: asegúrate de que tu MERGE es idempotente (repetirlo con los mismos datos da el mismo resultado). Los restarts del stream pueden aplicar la misma batch múltiples veces.
Trampa 1: "Al hacer UPDATE, se modifica el fichero Parquet en sitio"
❌ Los Parquet son inmutables. UPDATE genera nuevos ficheros y marca los viejos como removidos en el _delta_log.
Trampa 2: "Time travel funciona para siempre"
❌ Solo mientras los ficheros viejos existan. Después de VACUUM más allá del retention, ya no.
Trampa 3: "Puedo hacer VACUUM RETAIN 0 HOURS para limpiar todo"
⚠️ El sistema te lo impide por defecto (retention < 7 días bloqueado). Y hacerlo puede corromper lecturas concurrentes y romper time travel.
Trampa 4: "OPTIMIZE incluye VACUUM"
❌ Son operaciones distintas. OPTIMIZE consolida ficheros pequeños en grandes. VACUUM borra ficheros obsoletos.
Trampa 5: "V-Order hay que activarlo manualmente"
❌ Está activado por defecto en Fabric. Solo lo desactivas explícitamente si tienes motivos (staging heavy-write).
Trampa 6: "OptimizeWrite = OPTIMIZE"
❌ Son distintos: OptimizeWrite es preventivo, al escribir (activo por defecto); OPTIMIZE es un comando de mantenimiento post-hoc (manual/scheduled).
Trampa 7: "MERGE puede aceptar filas duplicadas en el source"
❌ Si hay filas duplicadas en el source que hacen match con la misma fila del target → MERGE falla. Deduplica el source antes.
Trampa 8: "En streaming, no necesito checkpointLocation"
❌ Es crítico para recuperación tras fallo y para exactly-once semantics.
Trampa 9: "Puedo compartir el mismo checkpoint entre streams"
❌ NO. Cada stream necesita su propio checkpoint location único.
Trampa 10: "Un stream desde Delta puede procesar UPDATEs del source"
❌ Por defecto solo procesa appends. Para UPDATE/DELETE del source, activa ignoreChanges o ignoreDeletes.
Trampa 11: "Schema evolution es automático siempre"
❌ Por defecto Delta enforcea schema (rechaza escrituras que no cumplen). Schema evolution hay que activarlo explícitamente con mergeSchema=true o config.
Trampa 12: "Particionar por columna de alta cardinalidad mejora performance"
❌ Genera miles de ficheros pequeños → peor performance. Usa baja cardinalidad.
Trampa 13: "DROP TABLE de una external borra los datos"
❌ Solo borra el metadata. Los datos permanecen en la ubicación especificada.
Trampa 14: "MERGE WITH SCHEMA EVOLUTION añade cualquier columna del source"
✅ Sí, si tiene columnas nuevas se añaden. Los tipos deben ser compatibles.
Trampa 15: "V-Order es propietario de Microsoft y no lo puede leer nadie más"
❌ V-Order es 100% Parquet compliant. Cualquier motor Parquet puede leerlo (Databricks, Snowflake, etc.). Solo que los motores de Fabric (VertiScan) lo aprovechan al máximo.