Cómo validar y alinear esquemas en ELT: paso a paso
Validar y alinear esquemas en ELT es esencial para evitar fallos de carga, problemas de calidad y regresiones cuando las fuentes cambian. Esta guía muestra cómo detectar diferencias de esquema, aplicar coerción/transformación y registrar cambios para que las cargas ELT sean robustas.
Requisitos previos
- Conocimientos básicos de SQL y familiaridad con un motor de procesamiento (ej.: Spark, Synapse, BigQuery).
- Acceso a un entorno ELT con capacidad de ejecutar transformaciones (ej.: cluster Spark o SQL warehouse).
- Ejemplos de archivos JSON/CSV o tablas de origen y una tabla de destino Delta/SQL.
Paso 1: Mapear esquemas esperados y reales
Antes de cargar, compara el esquema esperado (tabla destino) con el esquema real de la fuente. Esto permite identificar columnas ausentes, tipos incompatibles y columnas 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()
Paso 2: Detectar diferencias y reglas de coerción
Automatiza la comparación: crea lógica que detecte columnas ausentes, tipo diferente o mismatch de nullable. Define reglas de coerción (ej.: string->timestamp, int->long) y cómo tratar valores inválidos (null, default, error).
# Pseudocódigo en 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 diferencias
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)
Paso 3: Normalizar y convertir tipos en el ELT
Transforma la fuente al esquema destino aplicando cast, parsing y valores por omisión. Usa funciones robustas para evitar fallos (ej.: try_cast, to_timestamp con 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()
Paso 4: Registrar divergencias y crear informe de compatibilidad
Para mantenimiento, registra todas las divergencias encontradas con ejemplos y conteos. Esto ayuda a rastrear cambios de la fuente y a decidir si debes modificar el destino o aplicar transformaciones adicionales.
# 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')
Paso 5: Aplicar carga idempotente tras el alineamiento
Después de normalizar y validar, carga en la tabla destino usando una operación idempotente (ej.: MERGE o append controlado) para evitar duplicación e inconsistencias.
-- 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 el resultado
Confirma la validez verificando el esquema final, conteos y muestras. Ejecuta consultas para comparar conteos por partición y validaciones columna a columna (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;
Conclusión
Validar y alinear esquemas en ELT reduce fallos de carga y facilita la evolución de las pipelines. Próximos pasos: automatizar estas verificaciones con tests CI/CD y alertas; considera mantener un diccionario de datos y un proceso de versionado de esquema. Consejo: empieza por registrar pequeños informes de divergencia — son oro para diagnosticar cambios de la fuente.