Framework de validação e limpeza de dados com PySpark
Project description
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
- Executa pipeline antes/depois com comparativo
Requisitos
- Python 3.10+
- PySpark 3.4+
Instalacao
pip install -e .
Uso rapido
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
df_tratado, result_before, result_after = engine.run_dq_pipeline(
df=df,
table_name="CLIENTES",
output_dir="/mnt/dq_output", # opcional
)
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:
pip install -e ".[snowflake]"
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:
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.
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
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.5.5.tar.gz.
File metadata
- Download URL: dq_engine-0.5.5.tar.gz
- Upload date:
- Size: 41.4 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.11.0
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
f8d5e7d6813169115ef9286cf93c5903bb794fff1be73fd4161f7703b18acf57
|
|
| MD5 |
6fb3829ca9de034bd8064c30706f77c8
|
|
| BLAKE2b-256 |
0fbf5b9e8f956d5ad365ba65a186fb8828ddc8583371f9b04c2327e81b4beea2
|
File details
Details for the file dq_engine-0.5.5-py3-none-any.whl.
File metadata
- Download URL: dq_engine-0.5.5-py3-none-any.whl
- Upload date:
- Size: 46.8 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 |
960d84ad680d90c427744e1a9f253a18fae3cb92a0056fe5d125d8a13c6a7430
|
|
| MD5 |
a371cce881f73dfa8158a8c0c24e8031
|
|
| BLAKE2b-256 |
bdd0b46e1ed2effab7fab2036d19e803fa4ad05a44a550f0053a58892712b8cf
|