Ir al contenido
SparquetSparquet

sparquet_cola — biblioteca de calidad de datos

sparquet_cola es la capa de calidad de datos de Sparquet, extraída como una biblioteca que puedes usar por sí sola. Cola es la capa que pega calidad a tus DataFrames: checks de métrica estilo SODA Core, reglas SQL libres, verificación de esquema y el split válidas/inválidas (cuarentena) — todo sobre PySpark.

Depende solo de pyspark. Es intencional: puedes usarla en cualquier job Spark, notebook o task de Airflow sin arrastrar el resto del framework. Dentro de Sparquet es el motor detrás del bloque validations — los type de las reglas son idénticos, así que lo que aprendes aquí sirve igual en el JSON.

sparquet-cola es un paquete propio (repo). Instálalo standalone, o recíbelo automáticamente como dependencia de sparquet:

Terminal window
pip install sparquet-cola # standalone (solo pyspark)
pip install sparquet # el framework — trae sparquet-cola incluido
from sparquet_cola import Cola
cola = Cola()
Método Qué hace
cola.run(df, rules) Ejecuta todos los checks; devuelve una lista de CheckResult.
cola.split(df, rules) Divide el df en ColaSplit(valid, invalid) con los checks row-level.
cola.register(name, cls) Registra un check personalizado, disponible por type.
cola.available Lista los type registrados.

Cada regla es un dict con type + parámetros — el mismo shape del bloque validations.

results = cola.run(df, [
{"type": "row_count", "min": 1},
{"type": "missing_percent", "column": "cpf", "must_be": "< 5%", "warn": "= 0"},
])
for r in results:
print(r)

Un CheckResult trae: rule_type, passed, severity (pass/warn/fail), message, failed_count, metric_value y check_name.

Falla si alguna columna listada tiene un NULL. Row-level (entra en el split).

{"type": "not_null", "columns": ["id", "cpf"]}

Falla si la combinación de columnas no es única (clave compuesta). Row-level.

{"type": "unique", "columns": ["id_cesion", "numero_contrato"]}

Falla para valores fuera de [min, max] (inclusive; NULL pasa). Row-level.

{"type": "range", "column": "edad", "min": 0, "max": 150}

Falla para valores que no coinciden con el patrón (NULL cuenta como fallo). Row-level.

{"type": "regex", "column": "cpf", "pattern": "^[0-9]{11}$"}

Falla si el total queda fuera de [min, max]. Nivel de tabla.

{"type": "row_count", "min": 1, "max": 1000000}

Regla SQL libre sobre la temp view _validation_df. Dos modos:

  • query — invariante pass-when-true: devuelve un booleano; pasa cuando es true.
  • failed_rows — devuelve las filas malas (estilo SODA “failed rows”); falla si viene alguna. Con output, escribe esas filas en un destino.
{"type": "sql", "query": "SELECT COUNT(*) = 0 FROM _validation_df WHERE valor < 0",
"error_message": "hay valores negativos"}
{"type": "sql",
"failed_rows": "SELECT * FROM _validation_df WHERE valor < 0",
"output": {"format": "delta", "path": "dq.negativos", "mode": "overwrite"}}

Reglas de métrica — métrica + umbral (estilo SODA)

Sección titulada «Reglas de métrica — métrica + umbral (estilo SODA)»

Mide una métrica y la compara con un threshold, con niveles warn/fail.

Métricas: row_count, distinct_count, missing_count/missing_percent, duplicate_count/duplicate_percent, invalid_count/invalid_percent, min, max, avg, sum, stddev, freshness.

must_be es la condición de aprobación; warn (opcional) baja a aviso (no aborta).

DSL del threshold: >, <, >=, <=, =, !=, between X and Y, not between X and Y; sufijo % (porcentaje) y duración 1d/2h/30m (para freshness).

{"type": "row_count", "must_be": "> 0"}
{"type": "missing_percent", "name": "cpf completo",
"column": "cpf", "must_be": "< 1%", "warn": "= 0"}
{"type": "avg", "column": "valor", "must_be": "between 10 and 100"}
{"type": "freshness", "column": "actualizado_en", "must_be": "< 1d"}

Validez (para invalid_count/invalid_percent): una fila es inválida cuando está presente y viola la configuración de validez:

  • valid_values — conjunto aceptado · invalid_values — conjunto prohibido
  • valid_format — formato con nombre: email, uuid, cpf, cnpj, date, timestamp, phone, integer, decimal, url, ip, boolean, alphanumeric, credit_card, …
  • valid_regex — regex propia
  • valid_min / valid_max · valid_length / valid_min_length / valid_max_length
  • missing_values — valores tratados como ausentes además de NULL (para missing_*)
{"type": "invalid_percent", "column": "email",
"valid_format": "email", "must_be": "< 5%"}

Columnas obligatorias/prohibidas y tipos (data contract básico).

{"type": "schema",
"required_columns": ["id", "valor"],
"forbidden_columns": ["_debug"],
"column_types": {"id": "bigint", "valor": "double"}}

Los tipos coinciden sin distinguir mayúsculas y con alias (longbigint, integerint, strstring); decimal coincide por el tipo base, ignorando precisión/escala.

cola.split(df, rules) separa las filas usando los checks row-level (not_null, range, regex, unique, y las reglas de métrica missing/invalid). Una fila es inválida si viola cualquiera de ellos.

split = cola.split(df, [
{"type": "not_null", "columns": ["id"]},
{"type": "range", "column": "edad", "min": 0, "max": 150},
{"type": "invalid_count", "column": "email",
"valid_format": "email", "must_be": "= 0"},
])
split.valid.write.format("delta").mode("overwrite").save(".../silver_ok")
split.invalid.write.format("delta").mode("overwrite").save(".../silver_cuarentena")

Hereda de BaseCheck, implementa run() y, si es row-level, violation():

from pyspark.sql import functions as F
from sparquet_cola import Cola, CheckResult
from sparquet_cola.checks import BaseCheck
class NoFutureDateCheck(BaseCheck):
def run(self, df):
column = self.params["column"]
failed = df.filter(F.col(column) > F.current_date()).count()
if failed:
return CheckResult("no_future_date", False, f"{failed} fechas futuras", failed)
return CheckResult("no_future_date", True)
def violation(self, df):
return F.col(self.params["column"]) > F.current_date()
cola = Cola()
cola.register("no_future_date", NoFutureDateCheck)
cola.run(df, [{"type": "no_future_date", "column": "fecha_pedido"}])

El bloque validations de un pipeline JSON corre exactamente sobre este motor — mismos type, misma DSL de threshold, misma configuración de validez. El framework añade la persistencia del reporte, la política on_failure y la cuarentena row-level vía validations.outputs (valid/invalid). Consulta la referencia de validations y la guía de calidad de datos.