Skip to main content

Framework de validação e limpeza de dados com PySpark

Project description

dq-engine

Pacote Python dq-engine para 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

dq-engine/
├── 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, ColumnTreatmentContext
│       │   └── 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; NULL_RULE_TYPES
│       │   ├── data_clean.py    # Funções treat_* e apply_treatments: tratamento real dos dados
│       │   ├── outputs.py       # Builders do DataFrame de resultado
│       │   ├── result.py        # DataQualityResult: contrato de saída, comparativo, diagnóstico
│       │   ├── saver.py         # DqOutputSaver: persiste TXT e Excel (OneDrive/SharePoint)
│       │   ├── 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

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

# desenvolvimento local
pip install -e .

# pacote publicado no PyPI
pip install dq-engine

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 Top 3 valores inválidos por frequência com contagem, ex: `JJ (1240)
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

Tratamento de dados

O tratamento é aplicado diretamente via funções treat_* do módulo data_clean. Cada função recebe o nome da coluna e retorna uma Column Spark.

Via engine (aplica todas as colunas elegíveis)

# Aplica tratamentos em todas as colunas com regra mapeada
df_tratado, resumo = engine.apply_treatment(df, df_name="df_clientes")
print(resumo)
df_tratado.show(truncate=False)

Saída de exemplo:

Tratamentos aplicados via data_clean.apply_treatments():
  DAT_NASCIMENTO: date_null_pattern
  NOM_CLIENTE: string_null_pattern, test_value_pattern

Via data_clean diretamente

from dq_engine.execution.data_clean import treat_string, treat_date, treat_cpf

df_tratado = df.withColumns({
    "NOM_CLIENTE":     treat_string("NOM_CLIENTE"),
    "DAT_NASCIMENTO":  treat_date("DAT_NASCIMENTO"),
    "DOC_CPF_CLIENTE": treat_cpf("DOC_CPF_CLIENTE"),
})

Funções disponíveis:

Função Prefixos típicos Comportamento
treat_string NOM_, TPO_, NUM_, DOC_, END_, TEL_ Trim → upper → remove acentos → null → N/D
treat_cod COD_ Trim → remove acentos → remove não-alfanuméricos (sem UPPER); null → "0"
treat_dsc DSC_ Trim → upper → null → N/D (preserva espaços)
treat_idt IDT_ Domínio fechado: SIM / NAO / N/D
treat_date DAT_ Descarta datas inválidas → DateType; inválido → null
treat_val VAL_, QTD_ Cast double; null → 0.0
treat_uf *UF* Valida contra as 27 UFs; inválido → N/D
treat_cpf *CPF* Valida módulo 11 → 11 dígitos sem pontuação; inválido → N/D
treat_cnpj *CNPJ* Valida módulo 11 → 14 dígitos sem pontuação; inválido → N/D
treat_cep END_CEP* Remove pontuação → 8 dígitos; inválido → N/D
treat_placa *PLACA* Normatiza placa (padrão antigo e Mercosul); inválido → N/D
treat_chassi *CHASSI* Valida VIN ISO 3779; inválido → N/D
treat_telefone *TELEFONE*, *CELULAR* Valida DDD (Anatel) → (XX) XXXXX-XXXX; inválido → N/D

Pipeline completo (validação + tratamento + comparativo)

df_tratado, result_before, result_after, indicador, diagnostico = engine.run_dq_pipeline(
    df=df,
    table_name="tabela_clientes",
)

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.


Project details


Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

dq_engine-0.4.7.tar.gz (47.0 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

dq_engine-0.4.7-py3-none-any.whl (50.4 kB view details)

Uploaded Python 3

File details

Details for the file dq_engine-0.4.7.tar.gz.

File metadata

  • Download URL: dq_engine-0.4.7.tar.gz
  • Upload date:
  • Size: 47.0 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.2.0 CPython/3.11.0

File hashes

Hashes for dq_engine-0.4.7.tar.gz
Algorithm Hash digest
SHA256 e12cc0b011a5454d0bdcb7c03da19a6f56baa58b1120323a5f768ce53ebd973c
MD5 ede2683e4a95725c40da4c6859bb4a74
BLAKE2b-256 e7b31215df95c61daf3405ad55faee4b03972196fc73687af4ff25d29ea41b87

See more details on using hashes here.

File details

Details for the file dq_engine-0.4.7-py3-none-any.whl.

File metadata

  • Download URL: dq_engine-0.4.7-py3-none-any.whl
  • Upload date:
  • Size: 50.4 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.2.0 CPython/3.11.0

File hashes

Hashes for dq_engine-0.4.7-py3-none-any.whl
Algorithm Hash digest
SHA256 34645e9a0ae9b0a65379ff8e1da436744467ba88d2aa8fb4317e00f3accd215f
MD5 0fab1cae9213180e32364642eb61290b
BLAKE2b-256 d9581ca0601a4a927c332a3bd9fbeab768253bba92ec32011d43fa3825567919

See more details on using hashes here.

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