Framework de validação de qualidade de dados com PySpark
Project description
dq-engine
Framework de validação de qualidade de dados com PySpark.
Valida tabelas Spark com base em convenções de nomenclatura de colunas, regras declarativas por coluna, padrões de nulos e valores inválidos — gerando relatórios detalhados de inconsistências e sugestões de tratamento automático.
Sumário
- Objetivo
- Estrutura do projeto
- Arquitetura
- Instalação
- Como usar
- Interpretando o output
- Plano de tratamento
- Como adicionar novas regras
- Como adicionar novos prefixos
- Testes
Objetivo
O dq-engine valida tabelas Spark garantindo que os dados estejam de acordo com:
- Convenções de nomenclatura de colunas — prefixos como
COD_,DAT_,NOM_determinam automaticamente o conjunto de regras aplicável a cada coluna - Tipos de regras por coluna — trimming, upper case, padrões de nulo, regex, domínios enumerados
- Convenções de valores nulos — representações inválidas de nulo são detectadas e sinalizadas por tipo de coluna
- Valores de teste — detecta
TESTE,FAKE,***e similares em ambiente de produção
Estrutura do projeto
data-quality-local/
├── pyproject.toml
├── README.md
│
├── docs/
│ ├── architecture.md # Arquitetura, fluxo e decisões de design
│ └── rules.md # Regras de validação e rule sets
│
├── src/
│ └── dq_engine/
│ ├── __init__.py # Exportações públicas do pacote
│ │
│ ├── config/
│ │ ├── __init__.py
│ │ └── conventions.py # Fonte única de verdade: nulos, regras, prefixos, domínios
│ │
│ ├── planning/
│ │ ├── __init__.py
│ │ ├── models.py # Dataclasses: ValidationTarget, CompiledRule, ColumnTreatmentSpec
│ │ └── resolver.py # Resolução coluna → rule_set via prefix_mapping
│ │
│ ├── registry/
│ │ ├── __init__.py
│ │ └── registry.py # RuleRegistry: registro e lookup de builders de regras
│ │
│ ├── execution/
│ │ ├── __init__.py
│ │ ├── engine.py # DataQualityEngine: orquestração completa
│ │ ├── rule_compiler.py # SparkRuleCompiler: compila regras em Column expressions
│ │ ├── treatment_plan.py# Geração de código Spark de tratamento sugerido
│ │ ├── outputs.py # Builders do DataFrame de resultado
│ │ ├── result.py # DataQualityResult: contrato de saída
│ │ ├── validator.py # GenericValidator: validação pontual por coluna
│ │ └── exceptions.py # Hierarquia de exceções do framework
│ │
│ └── utils/
│ ├── __init__.py
│ └── logging.py # Logging estruturado
│
└── tests/
├── __init__.py
├── test_resolver.py # Testes do módulo planning/resolver
├── test_rule_compiler.py # Testes de cada builder de regra
├── test_engine.py # Testes de integração do engine
└── test_treatment_plan.py # Testes do gerador de código de tratamento
Arquitetura
DataFrame Spark
│
▼
planning/resolver.py → identifica colunas elegíveis por prefixo
│
▼
execution/rule_compiler.py → compila regras declarativas em Column expressions
│
▼
execution/engine.py → projeta flags de validação (1 select), agrega contagens (1 agg)
│
▼
execution/outputs.py → gera DataFrame de resultado por coluna
│
▼
DataQualityResult.output_df → 1 linha por coluna validada
Consulte docs/architecture.md para detalhes de cada componente e decisões de design.
Instalação
pip install -e .
Dependências:
- Python >= 3.10
- PySpark >= 3.4.0
Como usar
Validação de uma tabela
from pyspark.sql import SparkSession
from dq_engine import DataQualityEngine, CONVENTIONS
spark = SparkSession.builder.getOrCreate()
df = spark.table("meu_lakehouse.tabela_clientes")
engine = DataQualityEngine(spark=spark, conventions=CONVENTIONS)
result = engine.validate_table(
df=df,
table_name="tabela_clientes",
run_id="2024-01-15",
)
result.output_df.show(truncate=False)
Exemplo de output
+------------------+------------------+-------------------+----------+----------+------------+-------------+--------+
| table_name | column_name | invalid_value |total_rows|valid_rows|invalid_rows|quality_score| status |
+------------------+------------------+-------------------+----------+----------+------------+-------------+--------+
| tabela_clientes | NOM_CLIENTE | null | TESTE | 1000 | 985 | 15 | 0.985 | FAILED |
| tabela_clientes | COD_PRODUTO | | 1000 | 1000 | 0 | 1.0 | PASSED |
| tabela_clientes | DAT_NASCIMENTO | 0000-00-00 | 1000 | 998 | 2 | 0.998 | FAILED |
+------------------+------------------+-------------------+----------+----------+------------+-------------+--------+
Interpretando o output
| Coluna | Tipo | Descrição |
|---|---|---|
table_name |
string | Nome da tabela validada |
column_name |
string | Coluna com inconsistência |
invalid_value |
string | Amostra dos valores inválidos encontrados (até 50 distintos, separados por |) |
total_rows |
long | Total de linhas da tabela |
valid_rows |
long | Linhas que passaram em todas as regras da coluna |
invalid_rows |
long | Linhas que falharam em pelo menos uma regra (deduplificado) |
quality_score |
double | valid_rows / total_rows — entre 0.0 e 1.0 |
status |
string | PASSED se invalid_rows == 0, caso contrário FAILED |
Plano de tratamento
Além da validação, o engine gera automaticamente o código Spark de tratamento sugerido para corrigir as inconsistências encontradas.
# Opção 1: calcular contexto e gerar código em uma chamada
treatment_code = engine.render_treatment_code_from_df(
df=df,
df_name="df_clientes",
)
# Opção 2: calcular contexto separadamente (evita reprocessamento)
contexts = engine.get_column_treatment_contexts(df)
treatment_code = engine.render_treatment_code(df, contexts, df_name="df_clientes")
print(treatment_code)
Saída de exemplo:
df_clientes = df_clientes.select(
when(col("NOM_CLIENTE").isNull() | upper(trim(col("NOM_CLIENTE").cast("string"))).isin("NULL", "TESTE"),
lit("N/D")).otherwise(upper(trim(col("NOM_CLIENTE").cast("string")))).alias("NOM_CLIENTE"),
coalesce(col("COD_PRODUTO").cast("string"), lit("0")).alias("COD_PRODUTO"),
when(upper(col("DAT_NASCIMENTO").cast("string")).isin("0000-00-00"),
lit("1800-01-01 00:00:00").cast("timestamp")).otherwise(col("DAT_NASCIMENTO").cast("timestamp")).alias("DAT_NASCIMENTO"),
)
O código gerado pode ser copiado diretamente para o notebook de tratamento.
Como adicionar novas regras
1. Criar o builder em src/dq_engine/execution/rule_compiler.py
def _build_minha_regra(column_name: str, rule: dict, conventions: dict) -> CompiledRule:
condition = col(column_name).isNull() # sua condição de falha
return CompiledRule(
rule_type="minha_regra",
error_code="MINHA_REGRA_INVALIDA",
error_message="Descrição do erro",
condition=condition,
suggested_treatment_category="MANUAL_REVIEW",
suggested_treatment_expression=None,
manual_review_required=True,
treatment_note="Orientação de correção.",
)
2. Registrar no SparkRuleCompiler._register_default_rules()
def _register_default_rules(self) -> None:
builders = {
# ... builders existentes ...
"minha_regra": _build_minha_regra,
}
for rule_type, builder in builders.items():
self.registry.register(rule_type, builder)
3. Referenciar em um rule_set em src/dq_engine/config/conventions.py
"rule_sets": {
"meu_rule_set": [
{"type": "string_null_pattern"},
{"type": "minha_regra"},
],
}
4. Documentar em docs/rules.md
Como adicionar novos prefixos
Em src/dq_engine/config/conventions.py, adicione uma entrada em prefix_mapping:
"prefix_mapping": {
# ... prefixos existentes ...
"EMAIL": "email_field",
}
E defina o rule_set correspondente em rule_sets:
"email_field": [
{"type": "string_null_pattern"},
{"type": "regex", "pattern_ref": "email"},
],
Prefixos mais longos têm precedência automaticamente — nenhuma configuração extra necessária.
Testes
pytest tests/ -v
Os testes cobrem:
| Arquivo | Cobertura |
|---|---|
tests/test_resolver.py |
Resolução de prefixos, precedência, colunas sem match |
tests/test_rule_compiler.py |
Cada builder de regra individualmente |
tests/test_engine.py |
Fluxo completo de validação, quality_score, status |
tests/test_treatment_plan.py |
Geração de expressões e código de tratamento |