Pular para o conteúdo
SparquetSparquet

sparquet_cola — biblioteca de qualidade de dados

sparquet_cola é a camada de qualidade de dados do Sparquet, extraída como uma biblioteca que você pode usar sozinha. Cola é a camada que gruda qualidade nos seus DataFrames: checks de métrica estilo SODA Core, regras SQL livres, verificação de schema e o split válidas/inválidas (quarentena) — tudo em cima do PySpark.

Ela depende apenas de pyspark. Isso é proposital: dá para usá-la em qualquer job Spark, notebook ou task de Airflow sem trazer o resto do framework. Dentro do Sparquet ela é o motor por trás do bloco validations — os type das regras são idênticos, então o que você aprende aqui vale igual no JSON.

sparquet-cola é um pacote próprio (repo). Instale standalone, ou receba-o automaticamente como dependência do sparquet:

Terminal window
pip install sparquet-cola # standalone (só pyspark)
pip install sparquet # o framework — traz o sparquet-cola junto
from sparquet_cola import Cola
cola = Cola()
Método O que faz
cola.run(df, rules) Roda todos os checks; devolve uma lista de CheckResult.
cola.split(df, rules) Divide o df em ColaSplit(valid, invalid) pelos checks row-level.
cola.register(name, cls) Registra um check customizado, disponível por type.
cola.available Lista os type registrados.

Cada regra é um dict com type + parâmetros — o mesmo shape do bloco 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) # [WARN] 'cpf' check: missing_percent(cpf) = 2 passa must_be (< 5%) mas viola warn (= 0)

Um CheckResult traz: rule_type, passed, severity (pass/warn/fail), message, failed_count, metric_value e check_name.

Falha se qualquer coluna listada tiver algum NULL. Row-level (entra no split).

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

Falha se a combinação das colunas não for única (chave composta). Row-level.

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

Falha para valores fora de [min, max] (inclusive; NULL passa). Row-level.

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

Falha para valores que não casam com o padrão (NULL conta como falha). Row-level.

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

Falha se a contagem total ficar fora de [min, max]. Nível de tabela.

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

Regra SQL livre sobre a temp view _validation_df. Dois modos:

  • query — invariante pass-when-true: retorna um booleano; passa quando true.
  • failed_rows — retorna as linhas ruins (estilo SODA “failed rows”); falha se vier alguma. Com output, grava essas linhas num destino.
# invariante
{"type": "sql", "query": "SELECT COUNT(*) = 0 FROM _validation_df WHERE valor < 0",
"error_message": "há valores negativos"}
# failed rows (+ destino opcional)
{"type": "sql",
"failed_rows": "SELECT * FROM _validation_df WHERE valor < 0",
"output": {"format": "delta", "path": "dq.negativos", "mode": "overwrite"}}

Regras de métrica — métrica + threshold (estilo SODA)

Seção intitulada “Regras de métrica — métrica + threshold (estilo SODA)”

Mede uma métrica e compara com um threshold, com níveis 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 é a condição de aprovação; warn (opcional) rebaixa para aviso (não aborta).

DSL do threshold: >, <, >=, <=, =, !=, between X and Y, not between X and Y; sufixo % (percentual) e duração 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": "atualizado_em", "must_be": "< 1d"}

Validade (para invalid_count/invalid_percent): uma linha é inválida quando está presente e viola a configuração de validade:

  • valid_values — conjunto aceito · invalid_values — conjunto proibido
  • valid_format — formato nomeado: email, uuid, cpf, cnpj, date, timestamp, phone, integer, decimal, url, ip, boolean, alphanumeric, credit_card, …
  • valid_regex — regex própria
  • valid_min / valid_max · valid_length / valid_min_length / valid_max_length
  • missing_values — valores tratados como ausentes além de NULL (para missing_*)
{"type": "invalid_percent", "column": "email",
"valid_format": "email", "must_be": "< 5%"}

Colunas obrigatórias/proibidas e tipos (data contract básico).

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

Os tipos casam sem diferenciar maiúsculas e com aliases (longbigint, integerint, strstring); decimal casa pelo tipo base, ignorando precisão/escala.

cola.split(df, rules) separa as linhas usando os checks row-level (not_null, range, regex, unique, e as regras de métrica missing/invalid). Uma linha é inválida se viola qualquer um deles.

split = cola.split(df, [
{"type": "not_null", "columns": ["id"]},
{"type": "range", "column": "idade", "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_quarentena")

Herde de BaseCheck, implemente run() e, se for 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} datas 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": "dt_pedido"}])

O bloco validations de um pipeline JSON roda exatamente sobre este motor — mesmos type, mesma DSL de threshold, mesma configuração de validade. O framework acrescenta a persistência do relatório, a política on_failure e a quarentena row-level via validations.outputs (valid/invalid). Veja a referência de validations e o guia de qualidade de dados.