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
- 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
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
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file dq_engine-0.4.2.tar.gz.
File metadata
- Download URL: dq_engine-0.4.2.tar.gz
- Upload date:
- Size: 44.9 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.11.0
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
e5a59499a18ee2d8a10059ca7d6915261b0945dc2d2331c427ac6f84f3cdb1c9
|
|
| MD5 |
1c245f42b98fd8299142bf65d9c3a50c
|
|
| BLAKE2b-256 |
4e87ee5563715b87321787cc8573b3ab669638432e43937cd0582f53bd779c3c
|
File details
Details for the file dq_engine-0.4.2-py3-none-any.whl.
File metadata
- Download URL: dq_engine-0.4.2-py3-none-any.whl
- Upload date:
- Size: 48.1 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.11.0
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
5951d106051f57fab2498c79490a992680f5f7e43ac037665fc28c7a4f1d1273
|
|
| MD5 |
1a4a14701f56dc8371545112a4510cb1
|
|
| BLAKE2b-256 |
5df19501e6de7ed582366b29a33ef2b35c1587bd77341a622779e59cb417c500
|