Metadata-Version: 2.4
Name: dq-engine
Version: 0.7.0
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"
Provides-Extra: snowflake
Requires-Dist: snowflake-connector-python<4.0.0,>=3.6.0; extra == "snowflake"
Requires-Dist: pandas<3.0.0,>=2.0.0; extra == "snowflake"
Requires-Dist: pyyaml<7.0,>=6.0; extra == "snowflake"
Requires-Dist: python-dotenv<2.0.0,>=1.0.0; extra == "snowflake"
Requires-Dist: openpyxl<4.0.0,>=3.1.0; extra == "snowflake"

# dq-engine

Framework PySpark para validacao e tratamento de qualidade de dados.

## O que faz

- Resolve rule_set por nome de coluna
- Compila regras declarativas em expressoes Spark
- Gera score e status por coluna
- Aplica tratamento em lote (inclui colapso de espacos internos em texto DSC_)
- Executa pipeline antes/depois com comparativo
- Detecta avisos (informativos, nao alteram o score):
  - cardinalidade IDT_ fora do padrao;
  - similaridade categorica DSC_/TPO_ — mesma categoria com grafia diferente, por pontuacao/espacos (ex.: "S.A" vs "S/A");
  - nulos padronizados (`N/D`/`0`/`1800-01-01`) por coluna, marcando `[DROP?]` quando sao maioria.
- Registra em `dq_historico.xlsx` (abas Historico com `% Colunas Aprovadas`, Relatorio com `Causa`/`Acao Sugerida`, e Avisos)
- Suporta modo producao (valida o DataFrame ja tratado e marca situacao CORRIGIDO)

## Requisitos

- Python 3.10+
- PySpark 3.4+

## Instalacao

```bash
pip install -e .
```

## Uso rapido

```python
from pyspark.sql import SparkSession
from dq_engine import DataQualityEngine
from dq_engine.config.conventions import CONVENTIONS

spark = SparkSession.builder.getOrCreate()
engine = DataQualityEngine(spark=spark, conventions=CONVENTIONS)

df = spark.table("lakehouse.clientes")
result = engine.validate_table(df=df, table_name="CLIENTES", run_id="2026-07-02")

result.output_df.orderBy("column_name").show(truncate=False)
```

## Dois modos de uso

| Modo | Quando usar | O que faz |
|------|-------------|-----------|
| **Simulacao** (padrao) | Antes de levar o tratamento para o ETL | Valida o dado bruto, aplica o tratamento, valida de novo e compara ANTES x DEPOIS. Imprime o trecho de ETL (`treat_*`) para copiar e colar. |
| **Producao** (`producao=True`) | Depois que o ETL real ja aplicou o data clean | Valida **uma unica vez** o dado ja tratado (nao re-trata) e marca a situacao como CORRIGIDO, ou PENDENTE_FONTE se sobrar falha. |

### Simulacao

```python
df_tratado, result_before, result_after = engine.run_dq_pipeline(
    df=df,
    table_name="CLIENTES",
    output_dir="/mnt/dq_output",  # opcional: sem ele nada e gravado em arquivo
)
```

### Producao

```python
df_validado, result, _ = engine.run_dq_pipeline(
    df=df_ja_tratado,
    table_name="CLIENTES",
    output_dir="/mnt/dq_output",
    producao=True,
    ignore_columns_from_score=[],   # lista de colunas a tirar do score (veja abaixo)
)
```

### Ignorar colunas no score (`ignore_columns_from_score`)

Para uma particularidade ja conhecida e nao corrigivel agora (ex.: `IDT_X` que na verdade
deveria ser `TPO_X`), informe o nome exato da coluna:

```python
df_validado, result, _ = engine.run_dq_pipeline(
    df=df_ja_tratado,
    table_name="CLIENTES",
    output_dir="/mnt/dq_output",
    producao=True,
    ignore_columns_from_score=["IDT_SITUACAO", "IDT_SUBCATEGORIA"],
)
```

- So tem efeito com `producao=True` (em simulacao e ignorado).
- A coluna **continua sendo validada**; apenas fica fora do `quality_score`/indicador da
  tabela e do diagnostico de falha residual. Nao e manipular score, e priorizar.
- O estado real da coluna e mostrado no console/`.txt` (subsecao "COLUNAS MONITORADAS") e
  **gravado na aba `Avisos` do Excel** (`Tipo Aviso = COLUNA_MONITORADA`) para ajuste posterior.
- `[]` (ou omitir) = todas as colunas entram no score.

Via `main.py`: `python main.py --table ... --producao --ignore-columns-producao IDT_X`
(flag repetivel), ou no `.env`: `DQ_PRODUCAO=true` e `DQ_IGNORE_COLUMNS_PRODUCAO=IDT_X,IDT_Y`.

## Como ler o output

O output tem dois blocos de aviso:

1. **PONTO DE ATENCAO — nao reduz o score.** Um bloco unico, em subsecoes:
   - `IDT_` fora do padrao de dominio (mais de 2 valores distintos alem de `N/D`: provavel
     erro de prefixo, deveria ser `TPO_`/`DSC_`);
   - regras WARNING (ex.: ano suspeito em datas, `encoding_corruption`), **sempre com
     exemplos** dos valores que dispararam a regra;
   - categorias duplicadas por grafia (`S.A` x `S/A`);
   - colunas monitoradas (fora do score).
2. **FALHA RESIDUAL — reduz o score.** Colunas que seguem com falha apos o tratamento:
   corrigir na fonte.

Tudo isso vai tambem para o Excel (aba `Avisos`), com coluna, regra, categoria, quantidade e
mensagem (incluindo exemplos).

## Regras de tratamento de texto

Todo tratamento de string converte para **MAIUSCULAS**, remove acentos e trata nulos
(`N/D`), exceto:

- `COD_` — so `upper` + `trim`, preservando pontuacao (nulo vira `"0"`);
- colunas `DSC_` com `PATH`/`BUCKET`/`CAMINHO` (`treat_dsc_path`) — preservam o caminho
  e usam minusculas.

`DSC_` de baixa cardinalidade (<= `dsc_cardinality_threshold`) sao reclassificadas como
`dsc_category` e tratadas com `treat_string`; as demais usam `treat_dsc` (preserva espacos).

## Onde os arquivos sao salvos

Quando `output_dir` e informado, o `DqOutputSaver` grava:

```
{output_dir}/
├── table_output/
│   └── {TABELA}_{TIMESTAMP}.txt        # relatorio detalhado de cada execucao
└── output_consolidado/
    └── dq_historico_v2.xlsx            # historico acumulado (nome via historico_filename)
```

Abas do Excel: `Historico` (indicador/score por execucao), `Relatorio` (diagnostico por coluna),
`Avisos` (todos os avisos acima) e `Distribuicao` (opcional, `include_distribution=True`).

No `main.py` a configuracao e:

| Item | Flag / `.env` | Default |
|------|---------------|---------|
| Pasta de saida | `--output-dir` / `DQ_OUTPUT_DIR` | `\KIRIBATI.autoglass.com.br\Departamento$\Célula de Dados e Implantação\Engenharia\Outputs_Data_Quality` |
| Nome do Excel | `--historico-filename` / `DQ_HISTORICO_FILENAME` | `dq_historico_v2.xlsx` |
| Tabelas ja processadas | `--processed-file` / `DQ_PROCESSED_FILE` | `state/processed_tables.yaml` |
| Tabelas com falha | `--failed-file` / `DQ_FAILED_FILE` | `state/failed_tables.yaml` |

Se o Excel estiver aberto, o save avisa e nao grava o registro: feche e rode de novo.

## Leitura via Snowflake (opcional)

O engine sempre recebe um DataFrame Spark comum — a origem dos dados e desacoplada. Para descobrir e ler tabelas do Snowflake por schema, instale o extra `snowflake`:

```bash
pip install -e ".[snowflake]"
```

```python
from dq_engine.sources.snowflake_source import SnowflakeReader, discover_tables

reader = SnowflakeReader()  # credenciais via variaveis de ambiente (.env)
table_names = discover_tables(reader, schema_like="%NOME_DO_SCHEMA%")
df = reader.read(spark, table_names[0])

result = engine.validate_table(df=df, table_name=table_names[0], run_id="2026-07-14")
```

O script `main.py` automatiza o fluxo (descobre tabelas do schema -> roda dq-engine -> salva consolidado/historico) e pode ser agendado no Sonata:

```bash
python main.py --schema-like "%NOME_DO_SCHEMA%"
```

O padrao LIKE tambem pode vir do `.env` (variavel `TABLE_SCHEMA`), dispensando `--schema-like` em toda chamada.

Tambem e possivel rodar uma lista fixa de tabelas em vez de (ou combinada com) o schema, via `--table` (repetivel) ou pela variavel `TABLE_LIST` no `.env` (nomes completos `DATABASE.SCHEMA.TABELA` separados por virgula):

```bash
python main.py --table TRUSTED.MAX_APOLICE.TABELA_1 --table TRUSTED.MAX_APOLICE.TABELA_2
```

Tambem e possivel apontar um YAML com lista simples de tabelas (como `to_process.yaml`) usando `--tables-yaml` ou a variavel `TABLE_LIST_FILE` no `.env`:

```bash
python main.py --tables-yaml to_process.yaml
```

Sem `--schema-like`/`TABLE_SCHEMA`, o main.py roda so as tabelas informadas (`--table`, `TABLE_LIST` ou `--tables-yaml`), sem descoberta via `INFORMATION_SCHEMA`; com os dois informados, a lista filtra o resultado da descoberta.

Tabelas processadas com sucesso ficam registradas em `processed_tables.yaml` (arquivo temporario, desacoplado do dq-engine) e sao puladas nas proximas execucoes; use `--rerun-processed` para forcar tudo de novo.

## API principal

- DataQualityEngine.validate_table
- DataQualityEngine.apply_treatment
- DataQualityEngine.run_dq_pipeline
- DataQualityEngine.get_column_treatment_contexts
- DataQualityResult.table_summary
- DataQualityResult.failed_columns

## Configuracao

Arquivo central de convencoes:

- src/dq_engine/config/conventions.py

## Documentacao

- docs/architecture.md
- docs/rules.md
