SQLSaturday #1139 · Joinville · 11 Abr 2026
Qualidade
de Dados
no Databricks
Uma conversa sobre como definir, implementar e monitorar qualidade de dados em pipelines modernos com Databricks, desde os fundamentos até o hands-on com Great Expectations, DQX e frameworks personalizados.
<>
Patrocinadores
Case · 2024
JPMorgan Chase — US$350M de multa
01
Um dos maiores bancos do mundo, com bilhões de transações monitoradas por reguladores diariamente
02
Sistema de vigilância de trading recebia dados incompletos há anos, sem que ninguém tivesse detectado
03
Reguladores americanos identificaram a falha em 2024 durante auditoria de conformidade
04
Resultado: multa de US$350 milhões, revisão forçada de toda a infraestrutura de dados de trading
O que falhou: nenhum dado estava corrompido. Os dados chegavam, eram processados e pareciam corretos. Simplesmente faltavam campos que ninguém tinha definido como obrigatórios.
A lição
Dado incompleto não gera erro. Ele gera uma resposta errada que parece certa, até o regulador chegar.
Case · 2021
Zillow — quando o dado passa em tudo e ainda assim falha
01
Plataforma de estimativas imobiliárias que expandiu para compra e venda usando modelos preditivos próprios
02
Os dados eram tecnicamente válidos: sem nulos, sem erros, sem anomalias detectáveis no pipeline
03
O modelo não capturava contexto local, dinâmica de bairro, volatilidade pós-pandemia e sentimento de mercado
04
Resultado: ~US$500M de prejuízo, 25% dos funcionários demitidos, divisão encerrada
O que falhou: não foi o modelo de ML e não foi o pipeline. Os dados passariam em qualquer check. Faltava contexto que nenhuma validação técnica consegue capturar.
A diferença entre os dois casos
JPMorgan: dado faltando. Zillow: dado presente, mas insuficiente. Nos dois casos, ninguém viu vindo.
"Dados ruins não avisam
quando estão errados."
Agenda
O que vamos ver hoje
Teoria
  • Definições: ISO-25012 · DAMA DMBOK
  • Onde aplicar: Medallion Architecture
  • DQ vs Data Observability
  • Como planejar (5 etapas)
  • Expectativas · Contratos de Dados
  • Métricas de DQ · Anti-patterns
  • Ferramentas: panorama
Hands-on
  • Great Expectations
  • Framework personalizado
  • DQX
DQ
Definições
O que é
Qualidade
de Dados?
ISO-25012 · DAMA DMBOK
Dimensões de qualidade
ISO-25012 e DAMA DMBOK
ISO-25012 — norma internacional
Acurácia
Representa corretamente fatos do mundo real
Completude
Sem nulos em campos obrigatórios
Consistência
Sem contradições entre sistemas ou períodos
Atualidade
Dados dentro do período esperado de uso
Credibilidade
Origem confiável e processo de coleta adequado
DAMA DMBOK — acrescenta
Validade
Valores conformes às regras de negócio definidas
Integridade
Referências entre entidades funcionam, sem órfãos
Razoabilidade
Valores dentro dos limites esperados para o negócio
Unicidade
Chaves únicas sem duplicatas indevidas
Acurácia · Completude · Consistência · Atualidade aparecem em ambos os frameworks com definições equivalentes
02
Na prática
Data Quality
na Prática
Arquitetura
Onde aplicar DQ — Medallion Architecture
Bronze
Zona de pouso
Dados chegam como estão. O objetivo é ingerir com fidelidade, tratar depois.
  • Schema enforcement
  • NOT NULL em campos críticos
  • Volume mínimo esperado
  • Quarentena de inválidos
Depende da estratégia adotada
Silver
Limpeza e padronização
Dados limpos, tipados e normalizados, prontos para análise estruturada.
  • Limpeza e normalização
  • Tipagem correta
  • Deduplicação
  • Range checks e enums
  • Referential integrity
Gold
Consumo e produto
Dados orientados ao negócio com regras de domínio aplicadas.
  • Regras de negócio
  • Reconciliação de métricas
  • KPIs validados
  • SLA de freshness
  • Lakehouse Monitoring
Conceitos
Data Quality vs Data Observability
Data QualityData Observability
O que fazValida regras definidas explicitamenteDetecta anomalias e padrões inesperados
QuandoPontual, por execução do pipelineContínuo, monitora em tempo real
Exemploamount IS NOT NULL, enums válidosVolume caiu 40% sem motivo aparente
RespostaPass / Fail imediatoAlertas com contexto e histórico
FerramentasGX · DQX · Soda · dbt testsLakehouse Monitoring · Monte Carlo
Ambos são complementares. DQ define o contrato, Observability garante que ele está sendo cumprido.
Planejamento
Como planejar — 5 etapas
01
Definir expectativas
Schema, tipos, ranges, nulos, unicidade. Formalize antes de codar. Esse é o contrato.
02
Mapear pontos de validação
Bronze / Silver / Gold: cada camada tem seu conjunto de checks com criticidade definida.
03
Automatizar checks
Integre ao pipeline como etapa bloqueante ou com política de quarentena, não manual.
04
Monitorar continuamente
Freshness, volume, drift, métricas visíveis em dashboard, não apenas logs de erro.
05
Alertar e reagir
Defina runbook por tipo de falha. Monitoramento sem ação não resolve, resolve ansiedade.
03
Expectativas e Contratos
Formalizar
é o que
diferencia
Sem contrato, qualidade é uma conversa.
Contrato
Contrato de Dados
Um contrato de dados formaliza o acordo entre quem produz e quem consome um dado.
Produtor
Define schema, SLA e qualidade esperada. Notifica breaking changes com antecedência.
Consumidor
Declara dependências e expectativas. Aceita ou negocia o contrato explicitamente.
Versionamento
Contratos versionados, v1, v2, sem surpresa downstream quando o schema evolui.
Breaking changes
Renomear coluna, mudar tipo ou remover campo exige novo major com aviso prévio.
expectations:
  - title: "Uniqueness check"
    handler: uniqueness
    columns: [order_id]
    criticality: HIGH

  - title: "Timeliness SLA"
    handler: timeliness
    threshold: 1.0
    arrival_frequency: "0 9 * * *"
    stale_after: 3
    expected_within: 3

sla:
  freshness_hours: 24
  min_rows: 1000

owner: data-platform-team
version: 2.1.0
Papéis
Contrato de Dados na prática
PapelResponsabilidade no contrato
Engenheiro de DadosProdutor: define schema, SLA e publica o contrato
Tech Lead / DMAprova breaking changes e governa versionamento entre times
AnalistaConsumidor: declara dependências e expectativas sobre os dados
Cientista de DadosConsumidor: depende de garantias de qualidade para treinar e servir modelos
Data ManagerGarante que contratos existem e são respeitados entre times e domínios
Métricas
Métricas de Data Quality
% Completude
Por tabela e por coluna. Tendência importa mais que snapshot.
Freshness / Atraso
Tempo desde a última atualização vs SLA definido no contrato.
Volume: esperado vs real
Registros recebidos vs média histórica. Desvio acima de X% gera alerta.
Drift Rate
% de registros que falharam checks por execução. A tendência de degradação é o sinal mais valioso.
MTTR de incidentes
Tempo médio de resolução. Indicador de maturidade do processo.
Contratos violados
Número de data contracts com SLA quebrado por período.
Anti-patterns
Anti-patterns, o que não fazer
Validar só no final do pipeline
Erros no Gold custam mais do que erros no Bronze. Shift-left.
DQ como responsabilidade de um único time
Qualidade é responsabilidade do produtor, não do consumer ou de um time central.
Regras hardcoded no notebook
Sem versionamento, sem reuso, sem histórico. Use frameworks.
Monitoramento sem alertas e runbook
Dashboard bonito que ninguém olha não é monitoramento.
Falta de versionamento no contrato
Contrato muda silenciosamente. Downstream quebra sem aviso.
Ignorar dados históricos ruins
"Já estava assim" não é aceitação, é dívida técnica acumulando juros.
Ferramentas
Ferramentas, panorama de mercado
FerramentaTipoPonto forte
Great ExpectationsValidação (Python)Madura, rico ecossistema de expectativas, boa integração com Databricks
Soda CoreValidação (YAML)Agnóstica de plataforma, sintaxe simples e declarativa
dbt TestsValidação (SQL)Integrado ao workflow dbt, singular, not_null, custom SQL
DQXValidação (Databricks-native)API simples, Unity Catalog integrado, open source
DLT ExpectationsValidação (Databricks-native)Inline no pipeline, políticas warn / drop / fail, métricas no UI
Lakehouse MonitoringObservability (Databricks)Monitora drift e anomalias em tabelas Delta automaticamente
Monte CarloObservability (SaaS)Detecção de anomalias, lineage automático, alertas com contexto histórico
04
Prática
Hands-on
Great Expectations · Framework Personalizado · DQX
05
Hands-on · Parte 1
Great
Expectations
GX 1.x · Expectation Suites · Data Docs
Antes de rodar o código
Antes de rodar o código
Apresentacao GE.ipynb

Configurando o Great Expectations

In [ ]
%pip install great_expectations

%pip install -U "typing_extensions>=4.12"

dbutils.library.restartPython()
In [ ]
import great_expectations as gx
from great_expectations.exceptions import DataContextError

Data Context


É a configuração inicial, um espaço onde será salvo o contexto do projeto, como se fosse o init

In [ ]
# Cria o File Data Context na pasta atual
context = gx.get_context(mode="file", context_root_dir="/Workspace/Repos/arthurfr23@gmail.com/learn_databricks/src/great_expectations/gx")

Data_Source


Pode ser vários, desde pandas até spark, clouds como Snowflake etc.

É o tipo de source dos dados

In [ ]
data_source_name = "databricks.sql_saturday.pandas"
try:
    data_source = context.data_sources.get(data_source_name)
except Exception:
    data_source = context.data_sources.add_pandas(name=data_source_name)

Asset_name, batch request e expectation_suite_name


Asset_name são os dados. Utilizando o data_source, eles serão consultados nesse caminho


Batch request é o que será consultado no asset. Pode ser o asset completo (whole_df), uma quantidade limitada de linhas, um filtro de datas etc


Expectation_suite é uma coleção de expectativas para os dados

In [ ]
GOLD_TABLE_CONFIGS = [
    {"asset_name": "sql_saturday.gold.dim_parlamentares", "batch_def_name": "whole_df", "expectation_suite_name": "dim_parlamentares_suite"},
    {"asset_name": "sql_saturday.gold.dim_calendario", "batch_def_name": "whole_df", "expectation_suite_name": "dim_calendario_suite"},
    {"asset_name": "sql_saturday.gold.dim_fornecedores", "batch_def_name": "whole_df", "expectation_suite_name": "dim_fornecedores_suite"},
    {"asset_name": "sql_saturday.gold.fato_reembolso", "batch_def_name": "whole_df", "expectation_suite_name": "fato_reembolso_suite"},
]

Expectativas


As expectativas podem ser definidas de diferentes formas no Great Expectations.


Podemos defini-las manualmente, como fazemos abaixo, ou alterando diretamente nos yaml


O GE conta com diversas expectativas pré-configuradas, como o ExpectColumnValuesToNotBeNull, mas também podemos criar expectativas customizadas

In [ ]
def build_expectation_suite(suite_name: str, asset_name: str):
    """Constrói a ExpectationSuite com as expectativas específicas de cada tabela gold."""
    suite = gx.ExpectationSuite(name=suite_name)

    if "dim_parlamentares" in asset_name:
        # dim_parlamentares: not null, unicidade Id_Parlamentar, CPF 11 dígitos
        for col in ["Id_Parlamentar", "Nome_Parlamentar", "cpf", "UF", "Partido"]:
            suite.add_expectation(gx.expectations.ExpectColumnValuesToNotBeNull(column=col))
        suite.add_expectation(gx.expectations.ExpectColumnValuesToBeUnique(column="Id_Parlamentar"))
        suite.add_expectation(
            gx.expectations.ExpectColumnValueLengthsToBeBetween(column="cpf", min_value=11, max_value=11, strict_min=True, strict_max=True)
        )

    elif "dim_calendario" in asset_name:
        # dim_calendario: not null, unicidade, domínio Mes/Dia/Ano
        suite.add_expectation(gx.expectations.ExpectColumnValuesToBeUnique(column="Id_Calendario")
        )

    elif "dim_fornecedores" in asset_name:
        # dim_fornecedores: not null, unicidade Id e CNPJ_CPF
        for col in ["Id_Fornecedor", "Fornecedor", "CNPJ_CPF"]:
            suite.add_expectation(gx.expectations.ExpectColumnValuesToNotBeNull(column=col))
        suite.add_expectation(gx.expectations.ExpectColumnValuesToBeUnique(column="Id_Fornecedor"))
        suite.add_expectation(gx.expectations.ExpectColumnValuesToBeUnique(column="CNPJ_CPF"))

    elif "fato_reembolso" in asset_name:
        # fato_reembolso: not null FKs e chave, unicidade Id_Reembolso, valores >= 0
        for col in ["Id_Reembolso", "Id_Parlamentar", "Id_Calendario", "Id_Fornecedor"]:
            suite.add_expectation(gx.expectations.ExpectColumnValuesToNotBeNull(column=col))
        suite.add_expectation(gx.expectations.ExpectColumnValuesToBeUnique(column="Id_Reembolso"))
        for col in ["Valor_Documento", "Valor_Glosa", "Valor_Liquido", "Valor_Ressarcimento"]:
            suite.add_expectation(
                gx.expectations.ExpectColumnValuesToBeBetween(column=col, min_value=0, max_value=None, strict_min=True)
            )

    return suite

Linkando expectativas, assets e batchs

In [ ]
# Para cada tabela: get/create asset, get/create batch definition, get/create suíte (evita erro ao rodar de novo)
for config in GOLD_TABLE_CONFIGS:
    asset_name = config["asset_name"]
    batch_def_name = config["batch_def_name"]
    expectation_suite_name = config["expectation_suite_name"]

    try:
        data_asset = data_source.get_asset(asset_name)
    except Exception:
        data_asset = data_source.add_dataframe_asset(name=asset_name)

    try:
        data_asset.get_batch_definition(batch_def_name)
    except Exception:
        data_asset.add_batch_definition_whole_dataframe(batch_def_name)

    try:
        context.suites.get(name=expectation_suite_name)
        print(f"Suite '{expectation_suite_name}' já existe, pulando.")
    except Exception:
        suite = build_expectation_suite(expectation_suite_name, asset_name)
        context.suites.add(suite)
        print(f"Suite '{expectation_suite_name}' registrada para asset '{asset_name}'.")
  • Abrir o great_expectations.gx.expectations.dim_parlamentares_suite, onde estão as expectations
  • Abrir o great_expectations.gx.validation_definitions, onde estão as definitions de cada tabela
  • In [ ]
    for config in GOLD_TABLE_CONFIGS:
        validation_definition_name = config["expectation_suite_name"].replace("_suite", "_validation")
        batch_definition = (
            data_source.get_asset(config["asset_name"]).get_batch_definition(config["batch_def_name"])
        )
        suite = context.suites.get(name=config["expectation_suite_name"])
        vd = gx.ValidationDefinition(
            data=batch_definition,
            suite=suite,
            name=validation_definition_name,
        )
        try:
            context.validation_definitions.add(vd)
            print(f"Validation Definition '{validation_definition_name}' criada.")
        except DataContextError as e:
            if "already exists" in str(e):
                print(f"Validation Definition '{validation_definition_name}' já existe, pulando.")
            else:
                raise

    Configuração finalizada, agora vamos criar um notebook que rode as validações

    In [ ]
    import great_expectations as gx
    import great_expectations.expectations as gxe
    from great_expectations.validator.validator import Validator
    from great_expectations.checkpoint import UpdateDataDocsAction
    from great_expectations.exceptions import DataContextError
    from datetime import datetime, timezone
    from pyspark.sql import Row

    A classe abaixo recebe como parametros as o source_name, o suite, o asset, o batch, o expectation e o table_name.


    Com esses parametros ela busca no Data Context as informações, realiza os testes e retorna o resultado direto no output, em uma tabela e também no dashboard do GE.


    Após rodar a primeira parte do código 1x, não é necessário rodá-lo de novo. Agora podemos gerar alterações na expectation direto no yaml.

    In [ ]
    class GreatExpectationsExecute:
    
        def __init__ (self, 
                      data_source_name, 
                      suite_name, 
                      asset_name, 
                      batch_def_name, 
                      expectation_suite_name, 
                      table_name,
                      context_mode: str = "file",
                      context_root_dir: str = "/Workspace/Repos/arthurfr23@gmail.com/learn_databricks/src/great_expectations/gx"):
            
            self.data_source_name = data_source_name
            self.suite_name = suite_name
            self.asset_name = asset_name
            self.batch_def_name = batch_def_name
            self.expectation_suite_name = expectation_suite_name
            self.table_name = table_name
            self.context = gx.get_context(mode=context_mode, context_root_dir=context_root_dir)
    
        def data_source_name_check(self):
            try:
                self.data_source = self.context.data_sources.get(self.data_source_name)
            except Exception:
                print("Data Source não encontrado")
                raise
    
        def suite_name_check(self):
            try:
                self.suite = self.context.suites.get(self.suite_name)
            except Exception:
                print("Suite não encontrada")
                raise
    
        def asset_name_check(self):
            try:
                self.asset = self.data_source.get_asset(self.asset_name)
            except Exception:
                print("Asset name não encontrado")
                raise
    
        def batch_def_name_check(self):
            try:
                self.asset.get_batch_definition(self.batch_def_name)
            except Exception:
                print("Batch não definido")
                raise
    
        def expectation_suite_name_check(self):
            try:
                self.suite = self.context.suites.get(self.expectation_suite_name)
            except Exception:
                print("Expectation suite não encontrada")
                raise
    
        def table_name_check(self):
            try:
                self.table = spark.table(self.table_name)
            except Exception:
                print("Tabela não encontrada")
                raise
    
        def run(self, build_data_docs: bool = False, results_table: str = None):
            self.data_source_name_check()
            self.suite_name_check()
            self.asset_name_check()
            self.batch_def_name_check()
            self.expectation_suite_name_check()
            self.table_name_check()
            try:
                self._execute(results_table=results_table)
            finally:
                if build_data_docs:
                    self._geracao_html()
    
        def _execute(self, results_table: str = None):
            validation_definition_name = self.expectation_suite_name.replace("_suite", "_validation")
            vd = self.context.validation_definitions.get(validation_definition_name)
            df_spark = spark.read.table(self.asset_name)
            df = df_spark.toPandas()
            result = vd.run(batch_parameters={"dataframe": df})
            success = result.get("success", False) if isinstance(result, dict) else getattr(result, "success", False)
    #--------------------------------------------------------------------------------
            status = "OK" if success else "FALHOU"
            print(f"\n{'='*60}")
    
    #--------------------------------------------------------------------------------
            print(f"Asset: {self.asset_name} → {status}")
            print(f"{'='*60}")
    
            results_list = (
                result.get("results", []) if isinstance(result, dict)
                else getattr(result, "results", [])
            )
            for r in results_list:
                r_success = r.get("success", True) if isinstance(r, dict) else getattr(r, "success", True)
                expectation_type = (
                    r.get("expectation_config", {}).get("type", "")
                    if isinstance(r, dict)
                    else getattr(getattr(r, "expectation_config", None), "type", "")
                )
                kwargs = (
                    r.get("expectation_config", {}).get("kwargs", {})
                    if isinstance(r, dict)
                    else getattr(getattr(r, "expectation_config", None), "kwargs", {})
                )
                result_detail = (
                    r.get("result", {})
                    if isinstance(r, dict)
                    else getattr(r, "result", {})
                )
                icon = "✓" if r_success else "✗"
                print(f"  {icon} {expectation_type} | kwargs={kwargs} | result={result_detail}")
    
            if results_table:
                self._save_results(results_list, success, results_table)
    
            if not success:
                raise ValueError(
                    f"Validação GE falhou para {self.asset_name}. "
                    f"Ver resultado completo em result / Data Docs."
                )
            return status
    
        def _save_results(self, results_list: list, overall_success: bool, results_table: str):
            run_ts = datetime.now(timezone.utc)
            rows = []
            for r in results_list:
                r_success = r.get("success", True) if isinstance(r, dict) else getattr(r, "success", True)
                cfg = (
                    r.get("expectation_config", {}) if isinstance(r, dict)
                    else getattr(r, "expectation_config", {})
                )
                expectation_type = cfg.get("type", "") if isinstance(cfg, dict) else getattr(cfg, "type", "")
                kwargs = cfg.get("kwargs", {}) if isinstance(cfg, dict) else getattr(cfg, "kwargs", {})
                result_detail = (
                    r.get("result", {}) if isinstance(r, dict) else getattr(r, "result", {})
                )
                rows.append(Row(
                    run_timestamp=run_ts,
                    asset_name=self.asset_name,
                    expectation_suite=self.expectation_suite_name,
                    overall_success=overall_success,
                    expectation_type=expectation_type,
                    column=kwargs.get("column", None) if isinstance(kwargs, dict) else None,
                    success=bool(r_success),
                    unexpected_count=int(result_detail.get("unexpected_count", 0)) if isinstance(result_detail, dict) else None,
                    unexpected_percent=float(result_detail.get("unexpected_percent", 0.0)) if isinstance(result_detail, dict) else None,
                    kwargs_str=str(kwargs),
                    result_str=str(result_detail),
                ))
    
            df_results = spark.createDataFrame(rows)
            df_results.write.format("delta").mode("append").saveAsTable(results_table)
            print(f"\nResultados salvos em: {results_table}")
    
        def _geracao_html(self):
            docs = self.context.build_data_docs()
            if docs:
                for site_name, path in docs.items():
                    print(f"Data Docs ({site_name}): {path}")
                    print(f"Abra no navegador: file://{path}")
            else:
                print("Nenhum site de Data Docs configurado. Verifique great_expectations.yml.")
    In [ ]
    runner = GreatExpectationsExecute(
        data_source_name="databricks.sql_saturday.pandas",
        suite_name="dim_parlamentares_suite",
        asset_name="sql_saturday.gold.dim_parlamentares",
        batch_def_name="whole_df",
        expectation_suite_name="dim_parlamentares_suite",
        table_name="sql_saturday.gold.dim_parlamentares"
    )
    
    runner.run(build_data_docs=True, results_table="sql_saturday.monitoring.ge_results")
    In [ ]
    %sql
    SELECT *
    FROM sql_saturday.monitoring.ge_results

    Checar o output com valores de métricas e baixar o dashboard

    Como eu trabalharia:


  • Abrir pasta do tests, rodar o ge_configuration dentro do ge_utils,
  • criar o notebook .py
  • rodar os quality de cada tabela em um job
  • Escolheria entre usar o dashboard do próprio GE ou criar um personalizado com a tabela de controle no Databricks
  • Após rodar o código
    Após rodar o código
    dim_fornecedores_suite.json
    json
    {
      "expectations": [
        {
          "id": "8c635dcf-a459-4865-9a3d-7944153a3266",
          "kwargs": {
            "column": "Id_Fornecedor"
          },
          "meta": {},
          "severity": "critical",
          "type": "expect_column_values_to_not_be_null"
        },
        {
          "id": "43771305-4b0d-4061-9139-61f841eef3a6",
          "kwargs": {
            "column": "Fornecedor"
          },
          "meta": {},
          "severity": "critical",
          "type": "expect_column_values_to_not_be_null"
        },
        {
          "id": "248948ee-9ee8-490d-8f0a-6dae1e78b8fc",
          "kwargs": {
            "column": "CNPJ_CPF"
          },
          "meta": {},
          "severity": "critical",
          "type": "expect_column_values_to_not_be_null"
        },
        {
          "id": "79c9b68c-618f-4678-bd52-3fc7dbb5553f",
          "kwargs": {
            "column": "Id_Fornecedor"
          },
          "meta": {},
          "severity": "critical",
          "type": "expect_column_values_to_be_unique"
        },
        {
          "id": "0faa514d-d1d8-4e64-9b61-6bb8eb983523",
          "kwargs": {
            "column": "CNPJ_CPF"
          },
          "meta": {},
          "severity": "critical",
          "type": "expect_column_values_to_be_unique"
        }
      ],
      "id": "bf7e47cf-ae1d-4896-be49-85fbec12c54f",
      "meta": {
        "great_expectations_version": "1.16.0"
      },
      "name": "dim_fornecedores_suite",
      "notes": null
    }
    GE · Site 1
    GE · Site 1
    GE · Site 2
    GE · Site 2
    06
    Hands-on · Parte 2
    Frameworks
    Customizados
    PySpark · Delta Lake · Controle próprio
    Apresentacao Custom.ipynb

    **Checagens customizadas**

    A complexidade do GE pode assustar, porque sua configuração é um pouco complexa, existem processos que precisam ser feitos antes de rodar em produção etc.


    Por isso, em projetos mais rápidos, equipes que não conhecem o GE ou até mesmo por estratégia, é comum criarmos nosso próprio framework de testes

    In [ ]
    from pyspark.sql import functions as F
    from datetime import date

    Checagem de nulos

    In [ ]
    def check_nulls(table_name, col, criticality):
    
        test_date = date.today()
    
        test_type = 'check null'
    
        df = spark.read.table(f'sql_saturday.gold.{table_name}')
            
        nulos = df.select(F.count(F.when(F.col(col).isNull(), 1)).alias('num_nulos'))
    
        num_nulos = nulos.collect()[0][0]
    
        spark.sql(f"""
            INSERT INTO sql_saturday.custom_tests.tests_results (test_date, table_name, column_name, test_type, test_result, criticality)
            VALUES ('{test_date}', '{table_name}', '{col}', '{test_type}', {num_nulos}, '{criticality}')
        """)

    Checagem de unicidade

    In [ ]
    def check_unique(table_name, col, criticality):
    
        test_date = date.today()
    
        test_type = 'check unique'
    
        df = spark.read.table(f'sql_saturday.gold.{table_name}')
    
        df_duplicate = df.groupBy(col).count().filter("count > 1")
            
        total_duplicates = df_duplicate.agg(F.sum(col)).collect()[0][0]
    
        if total_duplicates is None:
            total_duplicates = 0
    
        spark.sql(f"""
            INSERT INTO sql_saturday.custom_tests.tests_results (test_date, table_name, column_name, test_type, test_result, criticality)
            VALUES ('{test_date}', '{table_name}', '{col}', '{test_type}', {total_duplicates}, '{criticality}')
        """)

    Tamanho de linhas

    In [ ]
    def check_lenght(table_name, col, criticality):
    
        test_date = date.today()
    
        test_type = 'check lenght'
    
        df = spark.read.table(f'sql_saturday.gold.{table_name}')
    
        df2 = df.withColumn('length_rows', F.length(df[col])).select('length_rows').filter('length_rows != 11').agg(F.sum('length_rows'))
    
        total_lenght = df2.collect()[0][0]
    
        if total_lenght is None:
            total_lenght = 0
    
        spark.sql(f"""
            INSERT INTO sql_saturday.custom_tests.tests_results (test_date, table_name, column_name, test_type, test_result, criticality)
            VALUES ('{test_date}', '{table_name}', '{col}', '{test_type}', {total_lenght}, '{criticality}')
        """)

    Valores positivos

    In [ ]
    def check_value_positive(table_name, col, criticality):
    
        test_date = date.today()
    
        test_type = 'check value positive'
    
        df = spark.read.table(f'sql_saturday.gold.{table_name}')
    
        total = df.select(F.count(F.when(F.col(col) <= 0, 1)).alias('non_positive')).collect()[0][0]
    
        if total is None:
            total = 0
    
        spark.sql(f"""
            INSERT INTO sql_saturday.custom_tests.tests_results (test_date, table_name, column_name, test_type, test_result, criticality)
            VALUES ('{test_date}', '{table_name}', '{col}', '{test_type}', {total}, '{criticality}')
        """)
    In [ ]
    check_unique('dim_calendario', 'Id_Calendario', 'critical')
    In [ ]
    %sql
    SELECT *
    FROM sql_saturday.custom_tests.tests_results

    Como eu trabalharia:

  • Criaria um arquivo.py com todas essas funções
  • Construiria um fluxo onde o .py lê um yaml, vê qual função, coluna, tabela e nível está configurado. Inclusive adicionando um filtro de data para não ler o dataset inteiro
  • Roda os testes em um job
  • Monitora via dashboard no Databricks


  • Por ser um framework customizado, podemos ir muito além, como criar tabelas de quarentena, dados validados etc

    Custom · DF1
    Custom · DF1
    Custom · DF2
    Custom · DF2
    Tabela de Controle
    Tabela de Controle
    Dashboard Custom
    Dashboard Custom
    07
    Hands-on · Parte 3
    DQX
    Databricks Labs · Profiling · Data Contracts · IA
    Apresentacao DQX.ipynb

    DQX

    Framework criado pelo próprio Databricks

  • Mais simples que o GE
  • Integrado com a plataforma
  • Tem funções com IA
  • Batch e streaming
  • Gera expectativas de forma automatica com o Profiling
  • Em um nível mais abrangente, tem duas funcionalidades principais:

  • Data Profiling, gerando estatísticas e expectativas automáticas com base nessas estatísticas
  • Data Validation, processo de validação de dados, com criação de tabelas de quarentena, por exemplo
  • In [ ]
    %pip install databricks-labs-dqx
    In [ ]
    %pip install 'databricks-labs-dqx[llm]'
    In [ ]
    %pip install 'databricks-labs-dqx[datacontract]'
    In [ ]
    dbutils.library.restartPython()

    1. Geração manual de testes

    In [ ]
    from databricks.labs.dqx import check_funcs
    from databricks.labs.dqx.rule import DQDatasetRule
    from databricks.labs.dqx.engine import DQEngine
    from databricks.sdk import WorkspaceClient
    In [ ]
    ws = WorkspaceClient()
    dq_engine = DQEngine(ws)
    In [ ]
    # Criação da expectation
    checks = [
      DQDatasetRule(
        name="id_calendar_month_is_unique",
        criticality="warn",
        check_func=check_funcs.is_unique,
        columns=["Id_Calendario, Mes"]
        )
      ]
    In [ ]
    # Execução do check
    input_df = spark.read.table('sql_saturday.gold.dim_calendario')
    result_df = dq_engine.apply_checks(input_df, checks)
    display(result_df)

    1. Geração automática por profiles

    Databricks lê as métricas do dataset e cria as próprias checagens

    In [ ]
    from databricks.labs.dqx.rule import DQRowRule
    from databricks.sdk import WorkspaceClient
    from databricks.labs.dqx.engine import DQEngine
    from databricks.labs.dqx.profiler.generator import DQGenerator
    from databricks.labs.dqx.profiler.profiler import DQProfiler
    from databricks.labs.dqx.config import WorkspaceFileChecksStorageConfig
    import logging
    import json
    import yaml
    In [ ]
    logging.basicConfig(level=logging.INFO)
    logger = logging.getLogger(__name__)
    In [ ]
    ws = WorkspaceClient()
    dq_engine = DQEngine(ws)
    generator = DQGenerator(ws)
    profiler = DQProfiler(ws)
    In [ ]
    input_df = spark.read.table("sql_saturday.gold.dim_calendario")
    In [ ]
    # Databricks faz a checagem e cria as regras, usando uma llm para detectar a PK
    ws = WorkspaceClient()
    profiler = DQProfiler(ws)
    summary_stats, profiles = profiler.profile(input_df, options={'llm_primary_key_detection': False})
    In [ ]
    generator = DQGenerator(ws)
    checks = generator.generate_dq_rules(profiles)
    In [ ]
    # Salva o profile e o check no workspace
    dq_engine.save_checks(checks, config=WorkspaceFileChecksStorageConfig(location="/Workspace/Repos/arthurfr23@gmail.com/learn_databricks/src/dqx_tests/dqx_apresentacao/dim_calendar_profile"))
    In [ ]
    checks = generator.generate_dq_rules(profiles)
    In [ ]
    # Carrega lista de checks a partir do arquivo YAML
    with open("/Workspace/Repos/arthurfr23@gmail.com/learn_databricks/src/dqx_tests/dqx_apresentacao/dim_calendar_profile", "r") as f:
        checks = yaml.safe_load(f)
    
    # Aplica os checks carregados aos dados
    result_df = dq_engine.apply_checks_by_metadata(input_df, checks)
    display(result_df)

    Alterar profile e rodar de novo a execução para ver erros

    In [ ]
    # Carrega lista de checks a partir do arquivo YAML
    with open("/Workspace/Repos/arthurfr23@gmail.com/learn_databricks/src/dqx_tests/dqx_apresentacao/dim_calendar_profile", "r") as f:
        checks = yaml.safe_load(f)
    
    # Aplica os checks carregados aos dados
    result_df = dq_engine.apply_checks_by_metadata(input_df, checks)
    display(result_df)

    3. Geração por IA

    In [ ]
    from databricks.labs.dqx.profiler.generator import DQGenerator
    from databricks.labs.dqx.config import InputConfig
    from databricks.sdk import WorkspaceClient
    In [ ]
    ws = WorkspaceClient()
    generator = DQGenerator(workspace_client=ws, spark=spark)
    
    user_input = """
    Id_calendario needs to be unique
    Ano need to be between 2000 and 2023
    Id_calendario needs to be between 1 and 365
    """
    
    checks = generator.generate_dq_rules_ai_assisted(
      user_input=user_input,
      input_config=InputConfig(location="sql_saturday.gold.dim_calendario")
    )
    
    print(checks)

    4. Geração por Data Contract

    In [ ]
    from databricks.labs.dqx.profiler.generator import DQGenerator
    from databricks.sdk import WorkspaceClient
    from databricks.labs.dqx.engine import DQEngine
    In [ ]
    ws = WorkspaceClient()
    engine = DQEngine(workspace_client=ws)
    generator = DQGenerator(workspace_client=ws, spark=spark)
    
    rules = generator.generate_rules_from_contract(
        contract_file="/Workspace/Repos/arthurfr23@gmail.com/learn_databricks/src/dqx_tests/dqx_apresentacao/dim_calendar_data_contract.yml"
    )
    In [ ]
    input_df = spark.read.table("sql_saturday.gold.dim_calendario")
    
    result_df = engine.apply_checks_by_metadata(
        input_df,
        rules
    )
    
    display(result_df)

    5. Métricas

    In [ ]
    from databricks.labs.dqx.engine import DQEngine
    from databricks.labs.dqx.metrics_observer import DQMetricsObserver
    from databricks.sdk import WorkspaceClient
    from databricks.labs.dqx.profiler.generator import DQGenerator
    from databricks.labs.dqx import check_funcs
    from databricks.labs.dqx.rule import DQDatasetRule
    from databricks.labs.dqx.config import InputConfig, OutputConfig
    In [ ]
    ws = WorkspaceClient()
    
    observer = DQMetricsObserver(name="dq_metrics")
    
    engine = DQEngine(WorkspaceClient(), observer=observer)
    
    generator = DQGenerator(workspace_client=ws, spark=spark)
    In [ ]
    # Lê o df e o Data Contract
    
    df = spark.read.table('sql_saturday.gold.dim_calendario')
    
    checks = generator.generate_rules_from_contract(
        contract_file="/Workspace/Repos/arthurfr23@gmail.com/learn_databricks/src/dqx_tests/dqx_apresentacao/dim_calendar_data_contract.yml"
    )
    In [ ]
    # Gera as métricas
    checked_df, observation = engine.apply_checks_by_metadata(df, checks)
    
    row_count = checked_df.count()
    
    
    metrics = observation.get
    print(f"Input row count: {metrics['input_row_count']}")
    print(f"Error row count: {metrics['error_row_count']}")
    print(f"Warning row count: {metrics['warning_row_count']}")
    print(f"Valid row count: {metrics['valid_row_count']}")
    In [ ]
    # Cria tabelas de validação, métricas e quarentenas para serem analisadas posteriormente
    input_config = InputConfig("sql_saturday.gold.dim_calendario")
    
    output_config = OutputConfig("sql_saturday.tests.valid_data")
    quarantine_config = OutputConfig("sql_saturday.tests.quarantine_data") 
    metrics_config = OutputConfig("sql_saturday.tests.metrics_data") 
    
    engine.apply_checks_by_metadata_and_save_in_table(
        checks=checks,
        input_config=input_config,
        output_config=output_config,
        quarantine_config=quarantine_config,
        metrics_config=metrics_config
    )
    In [ ]
    %sql
    SELECT *
    FROM sql_saturday.tests.valid_data
    In [ ]
    %sql
    SELECT *
    FROM sql_saturday.tests.metrics_data

    Olhar o monitoring

    DQX · DF1
    DQX · DF1
    DQX · DF2
    DQX · DF2
    dim_calendar_profile
    yaml
    - check:
        arguments:
          column: Id_Calendario
        function: is_not_null
      criticality: error
      name: Id_Calendario_is_null
    - check:
        arguments:
          column: Id_Calendario
          max_limit: 9364
          min_limit: 1
        function: is_in_range
      criticality: error
      name: Id_Calendario_isnt_in_range
    - check:
        arguments:
          column: Ano
        function: is_not_null
      criticality: error
      name: Ano_is_null
    - check:
        arguments:
          column: Ano
          max_limit: 2035
          min_limit: 2010
        function: is_in_range
      criticality: error
      name: Ano_isnt_in_range
    - check:
        arguments:
          column: Mes
        function: is_not_null
      criticality: error
      name: Mes_is_null
    - check:
        arguments:
          column: Mes
          max_limit: 12
          min_limit: 1
        function: is_in_range
      criticality: error
      name: Mes_isnt_in_range
    - check:
        arguments:
          column: Dia
        function: is_not_null
      criticality: error
      name: Dia_is_null
    - check:
        arguments:
          column: Dia
          max_limit: 31
          min_limit: 1
        function: is_in_range
      criticality: error
      name: Dia_isnt_in_range
    - check:
        arguments:
          column: Data_Completa
        function: is_not_null
      criticality: error
      name: Data_Completa_is_null
    - check:
        arguments:
          column: Data_Completa
          max_limit: '2035-08-21'
          min_limit: '2010-01-01'
        function: is_in_range
      criticality: error
      name: Data_Completa_isnt_in_range
    
    dim_calendar_data_contract.yml
    yaml
    kind: DataContract
    apiVersion: v3.0.2
    id: urn:datacontract:sql_saturday:gold_tables
    name: Regras das Tabelas Gold
    version: 2.0.0
    status: active
    domain: SQL Saturday
    dataProduct: sql_saturday_gold_rules
    
    schema:
      - name: dim_calendario
        physicalType: table
        properties:
          - name: Id_Calendario
            logicalType: string
            unique: true
          - name: Ano
            logicalType: string
            unique: true
    DQX · Monitoring Databricks
    DQX · Monitoring Databricks
    DQX · Print Dashboard
    DQX · Print Dashboard
    Conclusão
    Conclusão
    1
    DQ começa antes do pipeline, no contrato e na definição de expectativas. Sem contrato, não há qualidade.
    2
    Automatize validações e trate DQ como etapa bloqueante. Monitoramento sem ação não resolve, gera ansiedade.
    3
    Ferramenta certa + processo + ownership = dado confiável em produção. Nenhum dos três sozinho funciona.
    Arthur Ferreira Reis
    Arthur Ferreira Reis
    Engenheiro de Dados
    Databricks · Azure · Microsoft Fabric
    DP-700 DP-900 AZ-900 Databricks Lakehouse Databricks Analyst Student Ambassador
    LinkedIn QR
    LinkedIn