Ir al contenido
SparquetSparquet

Validaciones

Las validaciones miden el DataFrame después de cada transformación y antes de cualquier escritura. Nunca modifican los datos — esa separación es lo que hace confiable al reporte.

{
"validations": {
"on_failure": "warn",
"report": { "format": "csv", "path": "/dq/orders", "mode": "append" },
"rules": [
{ "type": "not_null", "columns": ["id", "customer_id"] },
{ "type": "unique", "columns": ["id"] },
{ "type": "range", "column": "quantity", "min": 1, "max": 10000 },
{ "type": "regex", "column": "email", "pattern": "^[^@\\s]+@[^@\\s]+\\.[a-z]{2,}$" },
{ "type": "row_count", "min": 1, "max": 5000000 },
{
"type": "sql",
"query": "SELECT COUNT(*) = 0 FROM _validation_df WHERE revenue < 0",
"error_message": "Negative revenue found"
}
]
}
}

Falla cuando cualquier columna listada contiene nulos. El conteo de fallas es el número de filas infractoras.

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

Falla cuando las columnas listadas no forman una clave única. Con varias columnas la unicidad es sobre la combinación, no sobre cada columna por separado.

{ "type": "unique", "columns": ["order_id", "line_number"] }

Acota una columna numérica o de fecha. Se requiere al menos uno de min / max; ambos son inclusivos.

{ "type": "range", "column": "quantity", "min": 1, "max": 10000 }

Los nulos pasan esta regla — combínala con not_null cuando la columna sea obligatoria.

Verifica una columna string contra un patrón.

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

Acota el tamaño del DataFrame — la guarda más barata contra una carga silenciosamente vacía.

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

min toma 0 por defecto; max es opcional.

Todo lo que las demás reglas no expresan. La query se ejecuta contra una vista temporal llamada _validation_df y debe devolver un único booleano: true significa que la regla pasó. La semántica es pasa-cuando-es-verdadero — escribe el invariante, no la violación.

{
"type": "sql",
"query": "SELECT COUNT(*) = 0 FROM _validation_df WHERE revenue < 0",
"error_message": "Negative revenue found"
}

Úsala para invariantes entre columnas (ends_at > starts_at), expectativas referenciales dentro del DataFrame, o reglas de negocio que solo tienen sentido en conjunto.

Dos modos. Además de la query booleana, la regla sql acepta failed_rows — una query que devuelve las filas infractoras (estilo “failed rows” de SODA). El check falla si regresa alguna, failed_count es su número, y un output opcional por regla escribe exactamente esas filas en un destino:

{
"type": "sql",
"failed_rows": "SELECT * FROM _validation_df WHERE revenue < 0",
"output": { "format": "delta", "path": "dq.negative_revenue", "mode": "overwrite" }
}

Usa query cuando solo necesitas pasa/falla; usa failed_rows cuando quieres ver y conservar las filas malas para depuración o una cuarentena dirigida.

Cada métrica es un type de regla: mide el DataFrame (o una columna) y compara el resultado con un umbral, al estilo SODA-Core. Las reglas de métrica llevan niveles de severidad warn y fail, así que una puede pasar y aun así levantar un aviso.

{ "type": "row_count", "must_be": "> 0" }
{ "type": "missing_percent", "name": "cpf rarely missing",
"column": "cpf", "must_be": "< 1%", "warn": "= 0" }
{ "type": "duplicate_count", "columns": ["id"], "must_be": "= 0" }
{ "type": "invalid_percent", "column": "email",
"valid_format": "email", "must_be": "< 5%" }
{ "type": "avg", "column": "amount", "must_be": "between 10 and 100" }
{ "type": "freshness", "column": "updated_at", "must_be": "< 1d" }

Métricas. row_count, distinct_count, missing_count / missing_percent, duplicate_count / duplicate_percent, invalid_count / invalid_percent, min, max, avg (mean), sum, stddev, y freshness (segundos desde el valor más reciente en column). Las métricas de conteo por columna toman column (o columns; las métricas de distinct/duplicate usan el conjunto completo).

DSL de umbral (en must_be, la condición de aprobación, y el warn opcional):

Forma Ejemplo
comparación > 0, < 5, >= 100, <= 10, = 0, != 0
rango between 10 and 20, not between 1 and 2
sufijo de porcentaje < 5% (el % es cosmético — el número es 5)
sufijo de duración < 1d, <= 2h, > 30m (convertido a segundos, para freshness)

Las unidades de duración son s m h d w. Un número solo significa igualdad (must_be: "0"= 0).

Validez de columna (para las métricas invalid_* — un valor cuenta como inválido cuando no está ausente y falla todas las reglas configuradas): valid_values, invalid_values, valid_format (regexes con nombre: email, uuid, cpf, cnpj, date, timestamp, phone, url, ip, integer, decimal, boolean, alphanumeric, credit_card, …), valid_regex, valid_min / valid_max, valid_min_length / valid_max_length / valid_length. Para las métricas missing_*, missing_values trata strings centinela extra (p. ej. ["", "N/A"]) como ausentes además de null.

Cualquier regla acepta targets: una lista de objetos donde cada uno se convierte en una regla propia. Todo lo que está fuera de targets es un default compartido; cada objetivo sobrescribe lo que necesita.

{ "type": "regex", "targets": [
{ "column": "cpf", "pattern": "^[0-9]{11}$" },
{ "column": "cnpj", "pattern": "^[0-9]{14}$" } ] }
{ "type": "range", "min": 1, "targets": [
{ "column": "id", "max": 1000 },
{ "column": "age", "max": 120 } ] }

Independiente es el punto: cada objetivo tiene su propio resultado, su propia fila en el reporte y su propio código de falla, así que el reporte dice cuál columna se rompió en lugar de dar un veredicto agregado. El range de arriba son dos reglas — range(id,1,1000) y range(age,1,120) — y una cuarentena puede acotarse a solo una de ellas.

La expansión ocurre al parsear la configuración, así que validations.rules ya está plano cuando las reglas se ejecutan — y el reporte, que empareja reglas con resultados por posición, se mantiene alineado.

Un objetivo no puede llevar type (una entrada es un único tipo de regla) ni targets anidado, y code / output van dentro de cada objetivo, no junto a la lista: de lo contrario todas las reglas expandidas compartirían un identificador o un destino. Cada una de esas formas se rechaza al parsear, en vez de producir un reporte ambiguo en silencio.

examples/08_validacao_multi_alvo.json ejecuta esto de punta a punta sobre CSV local: tres entradas de regla, siete filas en el reporte.

Asevera la forma del DataFrame — qué columnas deben existir, cuáles no, y sus tipos.

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

column_types acepta alias de tipo: longbigint, integerint, y una expectativa decimal(p,s) coincide con el tipo base decimal sin importar la precisión/escala.

Lo que hace el motor cuando al menos una regla falla.

Modo Pipeline Reporte Úsalo cuando
fail (por defecto) aborta antes de cualquier escritura no se escribe los datos malos nunca deben llegar al destino
warn continúa y escribe se escribe quieres visibilidad sin detener la carga
skip continúa y escribe se escribe las reglas son informativas por ahora

El modo se compara distinguiendo mayúsculas — escríbelo en minúsculas.

Las reglas de métrica agregan un nivel entre pasar y fallar. Una regla que supera must_be pero cruza su umbral opcional warn tiene severidad warn: se registra y reporta, pero cuenta como aprobado — un warn nunca aborta la ejecución, ni siquiera bajo on_failure: "fail". Solo una falla real (must_be violado, o cualquier otra regla que falle) lo hace.

Cada resultado lleva la severidad junto a los campos clásicos:

Campo Significado
severity pass, warn o fail — derivado de passed cuando una regla no lo define
metric_value el número que midió una regla de métrica (None para reglas que no producen métrica)
check_name el name opcional del check, para reportes legibles

Como un warn mantiene passed = True, un orquestador que filtra por not check.passed ve solo las fallas duras; inspecciona check.severity == "warn" para también sacar a la luz las advertencias.

report acepta una configuración de salida completa — cualquier formato, cualquier modo, incluso sus propias transformaciones. Una fila por regla:

Columna Significado
pipeline el name del pipeline
rule_type tipo de validador
check_name el name del check, cuando está definido (vacío en caso contrario)
severity pass, warn o fail
passed booleano (true para pass y warn)
failed_count filas infractoras (0 cuando pasó)
metric_value la métrica medida para las reglas check (null en caso contrario)
message detalle legible de la falla
validated_at timestamp del check
{
"report": {
"format": "delta",
"path": "quality.pipeline_report",
"mode": "append",
"partition_by": ["pipeline"]
}
}

Anexar cada ejecución a una sola tabla te da un historial de calidad que puedes graficar: tasa de fallas por regla, por pipeline, a lo largo del tiempo.

Cuarentena a nivel de fila (validations.outputs)

Sección titulada «Cuarentena a nivel de fila (validations.outputs)»

Donde report persiste métricas, validations.outputs enruta las filas mismas — aparte del/de los output(s) principal(es). Claves valid e invalid, cada una un destino de salida completo:

"validations": {
"on_failure": "warn",
"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" }
}
}

Una fila es inválida cuando viola cualquier regla a nivel de filanot_null, range, regex, unique, y las reglas de métrica missing_* / invalid_*. Las reglas agregadas (row_count, avg, freshness, duplicate_count, schema, sql booleano) describen la tabla, no filas individuales, así que no contribuyen a la división.

result = fw.run("orders.json")
for check in result.validation_results:
status = "ok" if check.passed else f"{check.failed_count} rows"
print(f"{check.rule_type:12} {status:>12} {check.message}")
if any(not check.passed for check in result.validation_results):
notify_data_owner(result.pipeline_name)
Intención Herramienta correcta
La fila no debe llegar al destino filter en transformations
Necesito saber cuántas filas malas llegaron una regla de validación
Los datos malos significan que la carga está mal una regla con on_failure: "fail"
Los datos malos se esperan pero deben rastrearse una regla con warn más un report
from sparquet.validation.base import BaseValidator, ValidationResult
import pyspark.sql.functions as F
class NoFutureDateValidator(BaseValidator):
def validate(self, df):
column = self.rule.params["column"]
failed = df.filter(F.col(column) > F.current_date()).count()
if failed:
return ValidationResult("no_future_date", False, f"{failed} future dates", failed)
return ValidationResult("no_future_date", True)
fw.register_validator("no_future_date", NoFutureDateValidator)
{ "type": "no_future_date", "column": "ordered_at" }

Consulta Extender.