(+351) 21 24 10006  ·  info@bconcepts.pt
Carnaxide, Lisboa

Cómo validar y alinear esquemas en ELT: paso a paso

João Barros 02 de September de 2026 5 min de lectura

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.