Pular para o conteúdo
SparquetSparquet

Validações

Validações medem o DataFrame após todas as transformações e antes de qualquer escrita. Elas nunca modificam os dados — essa separação é o que torna o relatório confiável.

{
"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"
}
]
}
}

Falha quando qualquer coluna listada contém nulos. A contagem de falhas é o número de linhas infratoras.

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

Falha quando as colunas listadas não formam uma chave única. Com várias colunas, a unicidade é sobre a combinação, não sobre cada coluna separadamente.

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

Limita uma coluna numérica ou de data. Ao menos um de min / max é obrigatório; ambos são inclusivos.

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

Nulos passam nesta regra — combine-a com not_null quando a coluna for obrigatória.

Verifica uma coluna string contra um padrão.

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

Limita o tamanho do DataFrame — a guarda mais barata contra uma carga silenciosamente vazia.

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

min tem default 0; max é opcional.

Qualquer coisa que as outras regras não expressam. A query executa contra uma temp view chamada _validation_df e precisa retornar um único booleano: true significa que a regra passou. A semântica é pass-when-true — escreva a invariante, não a violação.

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

Use-a para invariantes entre colunas (ends_at > starts_at), expectativas referenciais dentro do DataFrame, ou regras de negócio que só fazem sentido em conjunto.

Dois modos. Além da query booleana, a regra sql aceita failed_rows — uma query que retorna as linhas infratoras (estilo SODA “failed rows”). O check falha se alguma voltar, failed_count é o número delas, e um output opcional por regra escreve exatamente essas linhas para um sink:

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

Use query quando você só precisa de pass/fail; use failed_rows quando você quer ver e guardar as linhas ruins para depuração ou uma quarentena direcionada.

Cada métrica é um type de regra: ela mede o DataFrame (ou uma coluna) e compara o resultado a um threshold, no estilo SODA-Core. Regras de métrica carregam níveis de severidade warn e fail, então uma delas pode passar e ainda assim levantar um 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, e freshness (segundos desde o valor mais recente em column). Métricas de contagem por coluna aceitam column (ou columns; as métricas distinct/duplicate usam o conjunto inteiro).

DSL de threshold (em must_be, a condição de aprovação, e no warn opcional):

Forma Exemplo
comparação > 0, < 5, >= 100, <= 10, = 0, != 0
intervalo between 10 and 20, not between 1 and 2
sufixo de percentual < 5% (o % é cosmético — o número é 5)
sufixo de duração < 1d, <= 2h, > 30m (convertido para segundos, para freshness)

Unidades de duração são s m h d w. Um número puro significa igualdade (must_be: "0"= 0).

Validade de coluna (para métricas invalid_* — um valor conta como inválido quando é não-ausente e falha em toda regra configurada): valid_values, invalid_values, valid_format (regexes nomeadas: 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 métricas missing_*, missing_values trata strings sentinela extras (ex.: ["", "N/A"]) como ausentes além do null.

Qualquer regra aceita targets: uma lista de objetos em que cada um se torna uma regra própria. Tudo fora de targets é default compartilhado; cada alvo sobrescreve o que precisa.

{ "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 } ] }

Independente é o ponto: cada alvo tem seu próprio resultado, sua própria linha no relatório e seu próprio código de falha, então o relatório diz qual coluna quebrou em vez de dar um veredito agregado. O range acima são duas regras — range(id,1,1000) e range(age,1,120) — e uma quarentena pode ser escopada a apenas uma delas.

A expansão acontece no parse da configuração, então validations.rules já está achatado quando as regras rodam — e o relatório, que casa regras com resultados por posição, continua alinhado.

Um alvo não pode carregar type (uma entrada é um único tipo de regra) nem targets aninhado, e code / output ficam dentro de cada alvo, não ao lado da lista: senão todas as regras expandidas compartilhariam um identificador ou um destino. Cada uma dessas formas é recusada no parse, em vez de gerar um relatório ambíguo em silêncio.

O examples/08_validacao_multi_alvo.json roda isso de ponta a ponta em CSV local: três entradas de regra, sete linhas no relatório.

Afirma o formato do DataFrame — quais colunas devem existir, quais não devem, e seus tipos.

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

column_types aceita aliases de tipo: longbigint, integerint, e uma expectativa decimal(p,s) casa com o tipo base decimal independentemente de precisão/escala.

O que o motor faz quando ao menos uma regra falha.

Modo Pipeline Relatório Use quando
fail (default) aborta antes de qualquer escrita não escrito dados ruins nunca podem chegar ao destino
warn continua e escreve escrito você quer visibilidade sem interromper a carga
skip continua e escreve escrito as regras são informativas por ora

O modo é comparado de forma case-sensitive — escreva-o em minúsculas.

As regras de métrica adicionam um nível entre pass e fail. Uma regra que satisfaz must_be mas dispara seu threshold warn opcional tem severidade warn: ele é logado e reportado, mas conta como aprovado — um warn nunca aborta a execução, mesmo sob on_failure: "fail". Só uma falha real (must_be violado, ou qualquer outra regra falhando) o faz.

Cada resultado carrega a severidade junto aos campos clássicos:

Campo Significado
severity pass, warn ou fail — derivado de passed quando uma regra não a define
metric_value o número que uma regra de métrica mediu (None para regras que não produzem métrica)
check_name o name opcional do check, para relatórios legíveis

Como um warn mantém passed = True, um orquestrador filtrando por not check.passed vê apenas falhas duras; inspecione check.severity == "warn" para trazer os avisos também.

report aceita uma configuração de output completa — qualquer formato, qualquer modo, até suas próprias transformações. Uma linha por regra:

Coluna Significado
pipeline o name do pipeline
rule_type tipo do validator
check_name o name do check, quando definido (vazio caso contrário)
severity pass, warn ou fail
passed booleano (true para pass e warn)
failed_count linhas infratoras (0 quando passou)
metric_value a métrica medida para regras check (null caso contrário)
message detalhe de falha legível
validated_at timestamp do check
{
"report": {
"format": "delta",
"path": "quality.pipeline_report",
"mode": "append",
"partition_by": ["pipeline"]
}
}

Anexar cada execução numa única tabela lhe dá um histórico de qualidade que você pode plotar: taxa de falha por regra, por pipeline, ao longo do tempo.

Quarentena em nível de linha (validations.outputs)

Seção intitulada “Quarentena em nível de linha (validations.outputs)”

Onde report persiste métricas, validations.outputs roteia as próprias linhas — apartadas do(s) output(s) principal(is). Chaves valid e invalid, cada uma um sink de output 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" }
}
}

Uma linha é inválida quando viola qualquer regra em nível de linhanot_null, range, regex, unique, e as regras de métrica missing_* / invalid_*. Regras agregadas (row_count, avg, freshness, duplicate_count, schema, sql booleano) descrevem a tabela, não linhas individuais, então não contribuem para o split.

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)
Intenção Ferramenta certa
A linha não pode chegar ao destino filter em transformations
Preciso saber quantas linhas ruins chegaram uma regra de validação
Dados ruins significam que a carga está errada uma regra com on_failure: "fail"
Dados ruins são esperados mas precisam ser rastreados uma regra com warn mais um 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" }

Veja Estendendo.