Ir al contenido
SparquetSparquet

Calidad de datos en la práctica

Todo pipeline debería responder una pregunta antes de escribir: ¿estos datos son lo suficientemente buenos para publicarse? Las validaciones son la forma de que la respuesta pase a formar parte del archivo.

Para casi cualquier conjunto de datos, estas cuatro detectan la mayoría de los incidentes:

{
"validations": {
"on_failure": "warn",
"rules": [
{ "type": "row_count", "min": 1 },
{ "type": "not_null", "columns": ["id"] },
{ "type": "unique", "columns": ["id"] },
{ "type": "range", "column": "amount", "min": 0 }
]
}
}
  • row_count detecta la carga vacía silenciosa — un feed de origen que no llegó, un filtro que no coincidió con nada.
  • not_null sobre la clave detecta un join roto o una columna de origen renombrada.
  • unique sobre la clave detecta la lectura duplicada y el join que se abre en abanico.
  • range detecta errores de unidad e inversiones de signo.
Política El pipeline El reporte Cuándo es la adecuada
fail aborta antes de escribir no se escribe El destino alimenta algo que nunca debe ver filas malas
warn escribe de todos modos se escribe Quieres un historial y una señal, no una carga detenida
skip escribe de todos modos se escribe Las reglas todavía se están calibrando
result = fw.run("orders.json")
failed = [check for check in result.validation_results if not check.passed]
if failed:
alert(f"{result.pipeline_name}: {len(failed)} rules failed")
raise SystemExit(1) # your policy, in your orchestrator
{
"report": {
"format": "delta",
"path": "quality.pipeline_report",
"mode": "append",
"partition_by": ["pipeline"]
}
}

Cada ejecución agrega una fila por regla: pipeline, rule_type, check_name, severity, passed, failed_count, metric_value, message, validated_at. Apuntar todos los pipelines a la misma tabla te da un historial de calidad gratis, y tres gráficas que vale la pena tener:

  • tasa de fallo por regla a lo largo del tiempo — qué check realmente se gana su lugar
  • tendencia de failed_count por pipeline — la degradación antes de que se vuelva un incidente
  • reglas que nunca fallan — o los datos son sólidos o la regla está mal

Reglas para las cosas que realmente se rompen

Sección titulada «Reglas para las cosas que realmente se rompen»

Una vez que lo básico está en su lugar, las reglas valiosas son las reglas de dominio. sql las cubre: la consulta se ejecuta contra _validation_df y debe devolver un único booleano.

{
"type": "sql",
"query": "SELECT COUNT(*) = 0 FROM _validation_df WHERE ends_at < starts_at",
"error_message": "Contracts ending before they start"
}
{
"type": "sql",
"query": "SELECT ABS(SUM(debit) - SUM(credit)) < 0.01 FROM _validation_df",
"error_message": "Ledger does not balance"
}
{
"type": "sql",
"query": "SELECT COUNT(DISTINCT currency) = 1 FROM _validation_df",
"error_message": "Mixed currencies in a single settlement batch"
}

Esas tres valen más que una docena de checks de nulos, porque codifican lo que el negocio considera imposible.

Checks de métrica con una banda de advertencia

Sección titulada «Checks de métrica con una banda de advertencia»

Una regla de métrica mide una métrica y la compara con un umbral, al estilo SODA-Core — y puede advertir antes de fallar. Cada métrica es un type de regla: missing_percent, freshness, avg, row_count, … Esa banda de dos niveles es lo que permite que un pipeline se degrade de forma visible sin detenerse:

{ "type": "missing_percent", "name": "missing cpf",
"column": "cpf", "warn": "= 0", "must_be": "< 1%" }

Cero valores faltantes es el objetivo (warn), cualquier valor hasta el 1% se tolera (must_be), y a partir del 1% el check falla. Un warn se registra y aparece en el reporte, pero nunca aborta la ejecución, ni siquiera bajo on_failure: "fail" — así obtienes una señal temprana sin una falsa alarma.

La misma forma cubre la mayoría de las protecciones de métrica:

{ "type": "freshness", "column": "updated_at", "must_be": "< 1d" }
{ "type": "invalid_percent", "column": "email", "valid_format": "email", "must_be": "< 5%" }
{ "type": "avg", "column": "amount", "must_be": "between 10 and 100" }

Para las expectativas estructurales — las columnas y tipos de los que depende un consumidor aguas abajo — la regla schema falla rápido cuando un origen elimina silenciosamente una columna o cambia un tipo:

{ "type": "schema", "required_columns": ["id", "amount"], "column_types": { "id": "bigint", "amount": "double" } }

Consulta la referencia de validaciones para conocer cada métrica y el DSL completo de umbrales.

// removes the rows
{ "type": "filter", "condition": "id IS NOT NULL" }
// counts them, changes nothing
{ "type": "not_null", "columns": ["id"] }

Ambos tienen su lugar en la mayoría de los pipelines, y responden preguntas distintas. Un filtro mantiene el destino limpio; una regla te dice qué tan sucio estaba el origen. Solo la regla puede decirte que el feed se está degradando.

Delta e Iceberg lanzan multiple source rows matched cuando el DataFrame entrante tiene más de una fila por clave de merge. Una regla convierte esa caída en tiempo de ejecución en un reporte claro:

{
"validations": {
"on_failure": "fail",
"rules": [{ "type": "unique", "columns": ["order_id"] }]
},
"output": {
"format": "delta",
"path": "analytics.orders",
"mode": "merge",
"options": { "merge_keys": ["order_id"] }
}
}

La misma protección corresponde a cualquier escritura incremental con clave por un id de negocio.

Dos silvers: ¿cuarentena o salidas del pipeline?

Sección titulada «Dos silvers: ¿cuarentena o salidas del pipeline?»

Una necesidad recurrente: a partir de una tabla bronze, producir un silver válido y un silver inválido. Hay dos formas honestas de hacerlo, y la correcta depende de qué signifique “inválido”.

  • validations.outputs (cuarentena) — la división se deriva de las reglas: una fila es inválida cuando falla una regla a nivel de fila (not_null, range, regex, unique, o una regla de métrica missing_*/invalid_*). Una definición declarativa única de “inválido”, escrita junto a las métricas, que se ejecuta aparte de la salida principal — ni siquiera necesitas un output principal. Ideal cuando inválido significa literalmente “no pasó las reglas de calidad de datos”.
  • Dos outputs del pipeline, cada uno con su propia transformación filter — la división es lógica de negocio que escribes a mano. Ideal cuando “inválido” es una regla de dominio no relacionada con la calidad, o cuando las dos ramas necesitan formas distintas (columnas, joins, payloads).

Regla práctica: si de otro modo escribirías la misma condición dos veces — una como validación y otra como filtro — usa validations.outputs y deja que los checks sean la única fuente de verdad. Si las ramas divergen en significado o forma, usa las salidas del pipeline.

"validations": {
"rules": [
{ "type": "not_null", "columns": ["id"] },
{ "type": "invalid_count", "column": "email", "valid_format": "email", "must_be": "= 0" }
],
"outputs": {
"valid": { "format": "delta", "path": "silver.orders_ok", "mode": "overwrite" },
"invalid": { "format": "delta", "path": "silver.orders_quarantine", "mode": "overwrite" }
}
}

Para una cuarentena dirigida de las filas ofensoras de una sola regla, una regla sql con failed_rows más su propio output escribe exactamente esas filas. El motor completo es la librería independiente sparquet_cola — utilizable por sí sola en cualquier job de Spark.

  • row_count con un min en cada pipeline
  • not_null y unique sobre la clave de negocio
  • range sobre dinero, cantidades y fechas
  • regex sobre identificadores con una forma fija (documentos, códigos)
  • un sql para el invariante que tu dominio consideraría imposible
  • warn más un reporte persistido, salvo que el destino realmente no pueda tolerar filas malas
  • una regla unique en cualquier lugar donde un merge escriba