Cómo crear y usar una tabla Delta con histórico por día en Lakehouse
Este tutorial muestra cómo crear y usar una tabla Delta en Lakehouse que registra el historial diario de cambios (snapshot diario), permitiendo análisis temporales y auditoría. Es útil para conservar versiones por día sin recurrir al time travel continuo y para generar informes diarios consistentes.
Prerequisitos
- Acceso a un entorno Lakehouse (por ejemplo Microsoft Fabric Lakehouse) con permisos de escritura.
- Notebook con PySpark o endpoint SQL configurado para el Lakehouse.
- Un archivo CSV o streaming de ejemplo con datos de producto/stock para importar.
Paso 1: Pensar el modelo de historial diario
Antes de crear tablas, defina cómo almacenar el historial: una tabla Delta 'base' con los datos actuales y una tabla Delta 'snapshot_diario' que almacena una copia de las filas con una columna snapshot_date. Así se evita reescribir la base y se facilita el consumo por día.
Paso 2: Crear la tabla base Delta
Creé una tabla Delta donde quedan los registros actuales (esta tabla será actualizada por ETL/ingest). Mantenga la modelación simple: id, nome, stock, price, updated_at.
# PySpark example
spark.sql("CREATE TABLE IF NOT EXISTS lakehouse.produtos_base (
id STRING,
nome STRING,
stock INT,
price DOUBLE,
updated_at TIMESTAMP
) USING DELTA LOCATION 'abfss://@.dfs.core.windows.net/lakehouse/produtos_base'")
Paso 3: Poblar la tabla base con datos iniciales
Cargue un CSV de ejemplo en la tabla base. Si lo prefiere, haga un INSERT INTO o write DataFrame.
# PySpark example to load CSV and write to Delta
df = spark.read.option('header', 'true').csv('/mnt/data/produtos.csv')
from pyspark.sql.functions import col, current_timestamp
df2 = df.select(col('id'), col('nome'), col('stock').cast('int'), col('price').cast('double'))
.withColumn('updated_at', current_timestamp())
df2.write.format('delta').mode('overwrite').saveAsTable('lakehouse.produtos_base')
Paso 4: Crear la tabla de snapshots diarios
Creé una tabla Delta separada que tendrá una columna adicional snapshot_date de tipo date. Esta tabla acumulará una copia de los registros de la tabla base cada día.
spark.sql("CREATE TABLE IF NOT EXISTS lakehouse.produtos_snapshot_diario (
id STRING,
nome STRING,
stock INT,
price DOUBLE,
updated_at TIMESTAMP,
snapshot_date DATE
) USING DELTA LOCATION 'abfss://@.dfs.core.windows.net/lakehouse/produtos_snapshot_diario'")
Paso 5: Programar el proceso diario de snapshot (ejemplo manual)
Para probar, ejecute manualmente un job que inserte los datos de la tabla base en la tabla snapshot_diario con la fecha del día. En producción, programe esto con el scheduler del servicio (ej.: pipeline/flow).
# PySpark/SQL example to append daily snapshot
from pyspark.sql.functions import current_date
snapshot = spark.sql("SELECT id, nome, stock, price, updated_at FROM lakehouse.produtos_base")
snapshot.withColumn('snapshot_date', current_date()) \
.write.format('delta').mode('append').saveAsTable('lakehouse.produtos_snapshot_diario')
Paso 6: Manejar la desduplicación y solo registros cambiados
Para evitar duplicados innecesarios y reducir espacio, capture solo las filas que cambiaron desde el último snapshot. Haga un left-anti join entre la base y el último snapshot por id y valores relevantes.
# Exemplo: identificar alterações comparando com o último snapshot
last_snapshot = spark.sql("SELECT * FROM lakehouse.produtos_snapshot_diario WHERE snapshot_date = (SELECT max(snapshot_date) FROM lakehouse.produtos_snapshot_diario)")
base = spark.sql("SELECT id, nome, stock, price, updated_at FROM lakehouse.produtos_base")
changed = base.alias('b').join(last_snapshot.alias('s'), on='id', how='left') \
.filter("s.id IS NULL OR b.nome <> s.nome OR b.stock <> s.stock OR b.price <> s.price") \
.select('b.*')
changed.withColumn('snapshot_date', current_date()).write.format('delta').mode('append').saveAsTable('lakehouse.produtos_snapshot_diario')
Paso 7: Compactar y gestionar el crecimiento
Con acumulados diarios la tabla puede crecer. Periódicamente ejecute OPTIMIZE (si está disponible) y considere políticas de retención (ej.: conservar 365 días). Use VACUUM con cuidado conforme a la política del servicio.
# Exemplo SQL para limpar snapshots antigos (manter 365 dias)
-- Verifique primeiro as datas presentes
SELECT DISTINCT snapshot_date FROM lakehouse.produtos_snapshot_diario ORDER BY snapshot_date DESC;
-- Excluir linhas mais antigas (exemplo com condição WHERE)
DELETE FROM lakehouse.produtos_snapshot_diario WHERE snapshot_date < date_sub(current_date(), 365);
Verificar el resultado
Confirme que existen entradas en la tabla snapshot_diario para la fecha de hoy y que la base permanece intacta. Ejemplos de consultas para validar contenido y detectar duplicados.
-- Contagem por dia
SELECT snapshot_date, count(*) AS total FROM lakehouse.produtos_snapshot_diario GROUP BY snapshot_date ORDER BY snapshot_date DESC;
-- Ver datos de hoy
SELECT * FROM lakehouse.produtos_snapshot_diario WHERE snapshot_date = current_date() LIMIT 100;
-- Verifique se existem duplicados por id e dia
SELECT id, snapshot_date, count(*) FROM lakehouse.produtos_snapshot_diario GROUP BY id, snapshot_date HAVING count(*) > 1;
Conclusión
Ahora tiene una estrategia simple para crear y mantener un historial diario en Delta en el Lakehouse: una tabla base para datos actuales y una tabla snapshot_diario para análisis por día. Próximos pasos sugeridos: automatizar el job diario con un pipeline, añadir compresión/OPTIMIZE e implementar retención. Consejo: empiece por probar la lógica de desduplicación en conjuntos pequeños para evitar crecimiento descontrolado — ¿pretende guardar todo o solo los cambios importantes?