Como criar e usar uma tabela Delta com histórico por dia em Lakehouse
Este tutorial mostra como criar e usar uma tabela Delta em Lakehouse que regista o histórico diário de alterações (snapshot diário), permitindo análises temporais e auditoria. É útil para conservar versões por dia sem recorrer ao time travel contínuo e para gerar relatórios diários consistentes.
Pré-requisitos
- Acesso a um ambiente Lakehouse (por exemplo Microsoft Fabric Lakehouse) com permissões de escrita.
- Notebook com PySpark ou endpoint SQL configurado para o Lakehouse.
- Um ficheiro CSV ou streaming de exemplo com dados de produto/stock para importar.
Passo 1: Pensar o modelo de histórico diário
Antes de criar tabelas, defina como armazenar o histórico: uma tabela Delta 'base' com os dados actuais e uma tabela Delta 'snapshot_diario' que armazena uma cópia das linhas com uma coluna snapshot_date. Assim evita-se reescrever a base e facilita-se o consumo por dia.
Passo 2: Criar a tabela base Delta
Crie uma tabela Delta onde ficam os registos actuais (esta tabela será actualizada por ETL/ingest). Mantenha a modelação simples: 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'")
Passo 3: Popular a tabela base com dados iniciais
Carregue um CSV de exemplo para a tabela base. Se preferir, faça um INSERT INTO ou 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')
Passo 4: Criar a tabela de snapshots diários
Crie uma tabela Delta separada que terá uma coluna adicional snapshot_date do tipo date. Esta tabela vai acumular uma cópia dos registos da tabela base em cada dia.
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'")
Passo 5: Agendar o processo diário de snapshot (exemplo manual)
Para testar, execute manualmente um job que insere os dados da tabela base na tabela snapshot_diario com a data do dia. Em produção, agende isto com o scheduler do serviço (ex.: 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')
Passo 6: Lidando com deduplicação e apenas registos alterados
Para evitar duplicados desnecessários e reduzir espaço, capture apenas linhas que mudaram desde o último snapshot. Faça um left-anti join entre a base e o último snapshot por id e 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')
Passo 7: Compactar e gerir o crescimento
Com acumulados diários a tabela pode crescer. Periodicamente execute OPTIMIZE (se disponível) e considere políticas de retenção (ex.: conservar 365 dias). Use VACUUM com cuidado conforme a política do serviço.
# 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 o resultado
Confirme que existem entradas na tabela snapshot_diario para a data de hoje e que a base permanece intacta. Exemplos de consultas para validar conteúdo e 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 dados de hoje
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;
Conclusão
Agora tem uma estratégia simples para criar e manter um histórico diário em Delta no Lakehouse: uma tabela base para dados actuais e uma tabela snapshot_diario para análises por dia. Próximos passos sugeridos: automatizar o job diário com um pipeline, adicionar compressão/OPTIMIZE e implementar retenção. Dica: comece por testar a lógica de deduplicação em pequenos conjuntos para evitar crescimento descontrolado — pretende guardar tudo ou apenas as alterações importantes?