Metadata-Version: 2.4
Name: dq-engine
Version: 0.6.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`, Diagnostico 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)
```

## Pipeline completo

```python
df_tratado, result_before, result_after = engine.run_dq_pipeline(
    df=df,
    table_name="CLIENTES",
    output_dir="/mnt/dq_output",  # opcional
)
```

## Modo producao

Em producao o DataFrame ja chega tratado pelo ETL. Passe `producao=True` para validar uma
unica vez (sem reaplicar tratamento) e marcar a situacao como CORRIGIDO (ou PENDENTE_FONTE
quando ainda ha diagnostico residual):

```python
df_validado, result, _ = engine.run_dq_pipeline(
    df=df_ja_tratado,
    table_name="CLIENTES",
    output_dir="/mnt/dq_output",  # opcional
    producao=True,
)
```

Para uma particularidade já conhecida e não corrigível agora (ex.: coluna `IDT_X` que na
verdade deveria ser `TPO_X`), use `ignore_columns_from_score` — a coluna continua sendo
validada e aparece no output/aviso ("COLUNAS MONITORADAS"), mas fica de fora do
score/indicador da tabela (não é manipular o score, é priorização):

```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_X"],
)
```

Via `main.py`: `python main.py --table ... --producao --ignore-columns-producao IDT_X`
(repetível, ou `DQ_PRODUCAO=true` / `DQ_IGNORE_COLUMNS_PRODUCAO=IDT_X` no `.env`).

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