Como validar e alinhar esquemas em ELT: passo a passo
Validar e alinhar esquemas em ELT é essencial para evitar falhas de carga, problemas de qualidade e regressões quando as fontes mudam. Este guia mostra como detectar diferenças de esquema, aplicar coerção/transformação e registar alterações para que as cargas ELT sejam robustas.
Pré-requisitos
- Conhecimentos básicos de SQL e familiaridade com um motor de processamento (ex.: Spark, Synapse, BigQuery).
- Acesso a um ambiente ELT com capacidade de executar transformação (ex.: cluster Spark ou SQL warehouse).
- Exemplos de ficheiros JSON/CSV ou tabelas de origem e uma tabela de destino Delta/SQL.
Passo 1: Mapear esquemas esperados e reais
Antes de carregar, compare o esquema esperado (tabela destino) com o esquema real da fonte. Isto permite identificar colunas em falta, tipos incompatíveis e colunas extra.
-- Exemplo SQL para inspecionar esquema destino e fonte (Synapse/SQL genérico)
SELECT column_name, data_type
FROM information_schema.columns
WHERE table_name = 'destino_table'
ORDER BY ordinal_position;
-- Fonte: ler um ficheiro JSON com Spark para inferir esquema
val df = spark.read.option("multiline", true).json("/data/source/sample.json")
df.printSchema()
Passo 2: Detectar diferenças e regras de coerção
Automatize a comparação: crie lógica que detecte colunas em falta, tipo diferente ou mismatch de nullable. Defina regras de coerção (ex.: string->timestamp, int->long) e como tratar valores inválidos (null, default, erro).
# Pseudocódigo em PySpark para comparar esquemas
from pyspark.sql.types import StructType
def schema_to_dict(schema: StructType):
return {f.name: f.dataType.simpleString() for f in schema.fields}
src_schema = schema_to_dict(spark.read.json('/data/source/sample.json').schema)
dst_schema = schema_to_dict(spark.read.table('destino_table').schema)
# encontrar diferenças
missing = set(dst_schema) - set(src_schema)
diff_types = {col: (src_schema.get(col), dst_schema[col]) for col in src_schema if col in dst_schema and src_schema[col] != dst_schema[col]}
print('missing', missing)
print('diff_types', diff_types)
Passo 3: Normalizar e converter tipos no ELT
Transforme a fonte para o esquema destino aplicando cast, parsing e valores por omissão. Use funções robustas para evitar falhas (ex.: try_cast, to_timestamp com fallback).
# Exemplo PySpark: aplicar casts e colunas por omissão
from pyspark.sql.functions import col, lit, to_timestamp
df = spark.read.json('/data/source/sample.json')
# garantir coluna 'user_id' como long e 'created_at' como timestamp
df2 = df.withColumn('user_id', col('user_id').cast('long')) \
.withColumn('created_at', to_timestamp(col('created_at'), "yyyy-MM-dd'T'HH:mm:ss").alias('created_at'))
# adicionar colunas em falta com valor por omissão
for c in ['status', 'country']:
if c not in df2.columns:
df2 = df2.withColumn(c, lit(None).cast('string'))
df2.printSchema()
Passo 4: Registar divergências e criar relatório de compatibilidade
Para manutenção, regista todas as divergências encontradas com exemplos e contagens. Isto ajuda a rastrear alterações de fonte e a decidir se tens de alterar destino ou aplicar transformações adicionais.
# Exemplo: contar valores que falharam no cast
bad_userid = df.filter(col('user_id').cast('long').isNull() & col('user_id').isNotNull()).count()
print(f'valores user_id com cast inválido: {bad_userid}')
# Escrever relatório simples numa tabela de auditoria
report = spark.createDataFrame([('missing_cols', ','.join(missing)), ('bad_userid', str(bad_userid))], ['issue', 'detail'])
report.write.mode('append').saveAsTable('audit_schema_issues')
Passo 5: Aplicar carregamento idempotente após alinhamento
Depois de normalizar e validar, carrega para a tabela destino usando uma operação idempotente (ex.: MERGE ou append controlado) para evitar duplicação e inconsistências.
-- Exemplo MERGE em SQL (Delta Lake)
MERGE INTO destino_table AS d
USING staging_table AS s
ON d.id = s.id
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *;
-- Alternativa: overwrite particionado apenas para partição específica
INSERT OVERWRITE destino_table PARTITION(dt='2026-09-01') SELECT * FROM staging_table WHERE dt='2026-09-01';
Verificar o resultado
Confirma a validade verificando esquema final, contagens e amostras. Executa queries para comparar contagens por partição e validações coluna-a-coluna (tipos, nulls, formato de timestamps).
-- Verificar esquema final
SELECT column_name, data_type, is_nullable FROM information_schema.columns WHERE table_name='destino_table';
-- Contagens por partição e amostra
SELECT dt, COUNT(*) FROM destino_table GROUP BY dt ORDER BY dt DESC LIMIT 10;
-- Amostra de linhas com valores inválidos
SELECT * FROM destino_table WHERE TRY_CAST(user_id AS BIGINT) IS NULL AND user_id IS NOT NULL LIMIT 10;
Conclusão
Validar e alinhar esquemas em ELT reduz falhas de carga e facilita evolução das pipelines. Próximos passos: automatizar estas verificações com testes CI/CD e alertas; considera manter um dicionário de dados e um processo de versionamento de esquema. Dica: começa por registar pequenos relatórios de divergência — são ouro para diagnosticar alterações de fonte.