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.
Instalación e import
Sección titulada «Instalación e import»sparquet-cola es un paquete propio (repo). Instálalo standalone, o recíbelo automáticamente como dependencia de sparquet:
pip install sparquet-cola # standalone (solo pyspark)pip install sparquet # el framework — trae sparquet-cola incluidofrom sparquet_cola import Colacola = Cola()La API Cola
Sección titulada «La API 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.
Checks disponibles
Sección titulada «Checks disponibles»not_null
Sección titulada «not_null»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}$"}row_count
Sección titulada «row_count»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 estrue.failed_rows— devuelve las filas malas (estilo SODA “failed rows”); falla si viene alguna. Conoutput, 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 prohibidovalid_format— formato con nombre:email,uuid,cpf,cnpj,date,timestamp,phone,integer,decimal,url,ip,boolean,alphanumeric,credit_card, …valid_regex— regex propiavalid_min/valid_max·valid_length/valid_min_length/valid_max_lengthmissing_values— valores tratados como ausentes además de NULL (paramissing_*)
{"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 (long→bigint, integer→int, str→string); decimal coincide por el tipo base, ignorando precisión/escala.
Split válidas/inválidas (cuarentena)
Sección titulada «Split válidas/inválidas (cuarentena)»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")Check personalizado
Sección titulada «Check personalizado»Hereda de BaseCheck, implementa run() y, si es row-level, violation():
from pyspark.sql import functions as Ffrom sparquet_cola import Cola, CheckResultfrom 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"}])Dentro del framework
Sección titulada «Dentro del framework»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.