Skip to main content

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


Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

dq_engine-0.5.5.tar.gz (41.4 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

dq_engine-0.5.5-py3-none-any.whl (46.8 kB view details)

Uploaded Python 3

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

Hashes for dq_engine-0.5.5.tar.gz
Algorithm Hash digest
SHA256 f8d5e7d6813169115ef9286cf93c5903bb794fff1be73fd4161f7703b18acf57
MD5 6fb3829ca9de034bd8064c30706f77c8
BLAKE2b-256 0fbf5b9e8f956d5ad365ba65a186fb8828ddc8583371f9b04c2327e81b4beea2

See more details on using hashes here.

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

Hashes for dq_engine-0.5.5-py3-none-any.whl
Algorithm Hash digest
SHA256 960d84ad680d90c427744e1a9f253a18fae3cb92a0056fe5d125d8a13c6a7430
MD5 a371cce881f73dfa8158a8c0c24e8031
BLAKE2b-256 bdd0b46e1ed2effab7fab2036d19e803fa4ad05a44a550f0053a58892712b8cf

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page