Skip to main content

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

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

Project details


Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page