Metadata-Version: 2.4
Name: dq-engine
Version: 0.4.3
Summary: Framework de validação e limpeza de dados com PySpark
Author-email: Bruna Cataldo <bruna.cataldo@autoglass.com.br>
Project-URL: Homepage, https://bitbucket.org/ced_engenharia/data_quality
Requires-Python: >=3.10
Description-Content-Type: text/markdown
Requires-Dist: pyspark<4.0.0,>=3.4.0
Provides-Extra: dev
Requires-Dist: pytest<9.0.0,>=8.0.0; extra == "dev"
Requires-Dist: pytest-mock<4.0.0,>=3.12.0; extra == "dev"
Provides-Extra: notebook
Requires-Dist: openpyxl<4.0.0,>=3.1.0; extra == "notebook"
Requires-Dist: pandas<3.0.0,>=2.0.0; extra == "notebook"

# 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](#objetivo)
- [Estrutura do projeto](#estrutura-do-projeto)
- [Arquitetura](#arquitetura)
- [Instalação](#instalação)
- [Como usar](#como-usar)
- [Interpretando o output](#interpretando-o-output)
- [Plano de tratamento](#plano-de-tratamento)
- [Como adicionar novas regras](#como-adicionar-novas-regras)
- [Como adicionar novos prefixos](#como-adicionar-novos-prefixos)
- [Testes](#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`](docs/architecture.md) para detalhes de cada componente e decisões de design.

---

## Instalação

```bash
# 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

```python
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) | XX (89)` |
| `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)

```python
# 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

```python
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)

```python
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`

```python
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()`

```python
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`

```python
"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`:

```python
"prefix_mapping": {
    # ... prefixos existentes ...
    "EMAIL": "email_field",
}
```

E defina o rule_set correspondente em `rule_sets`:

```python
"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.

---
