A Jornada Databricks será uma série de artigos em que compartilho o que venho aprendendo no meu dia a dia com a ferramenta. O objetivo é consolidar meus aprendizados e experiências enquanto construo uma plataforma de dados utilizando o ecossistema da Databricks.
Como este é o primeiro artigo, preciso contextualizar algumas escolhas que fiz e que se distanciam um pouco do fluxo padrão do Spark Declarative Pipelines ou Lakeflow Pipelines, como é chamado na Databricks.
Background
Atualmente, nossa plataforma de dados é construída em cima de um projeto dbt, utilizando Apache Iceberg como formato de tabela e Trino como query engine. Surgiu uma necessidade de negócio para que determinados processamentos fossem atualizados em janelas de 1 hora. Para o cenário tradicional de D-1, a arquitetura dbt + Trino sempre funcionou muito bem. Contudo, ao migrar para atualizações horárias, o desafio aumentou consideravelmente, pois o custo computacional de identificar o que mudou para aplicar uma carga incremental tornou-se maior do que simplesmente reprocessar a tabela inteira (full refresh).
Acontece que reprocessar tabelas inteiras a cada hora gera um custo gigantesco de computação e armazenamento, além de fazer com que o uso de formatos com suporte a time travel traga desvantagens dado que os snapshots e versões antigas do Iceberg acumulam-se rapidamente como lixo, demandando um esforço operacional contínuo de manutenção.
Atento a essa dor, eu já vinha acompanhando o SDP (Spark Declarative Pipelines), lançado no Spark 4.1 e derivado do DLT (Delta Live Tables, atual Lakeflow Pipelines). Ele parecia uma excelente alternativa por se basear no conceito de View Materializadas. No entanto, o SDP open source é apenas um Lakeflow Pipelines simplificado. Ele não possui suporte ao IVM (Incremental View Maintenance), recurso essencial para a nossa necessidade de atualização incremental eficiente.
Por outros motivos estratégicos, a empresa onde trabalho optou por contratar a Databricks, o que acabou resolvendo essa e outras dores. Essa necessidade de reprocessamento horário não foi o principal causador da mudança, o motivador real foi o tamanho reduzido do nosso time (apenas 3 pessoas). Estava inviável manter uma infraestrutura 100% open source e, ao mesmo tempo, evoluir a plataforma. Para se ter uma ideia, mantínhamos internamente serviços como Trino, Trino Gateway, Starrocks, Cube, Superset, Lakekeeper, Airflow, OpenMetadata, dbt, Meltano e conectores do Kafka (Debezium e Iceberg Sink).
Primeiras dores
A partir de agora, falarei especificamente do ecossistema Databricks, sem me aprofundar nas limitações da versão open source do SDP por ser bastante reduzida em comparação ao Lakeflow Pipelines.
Logo de início, busquei entender como seria a esteira de desenvolvimento e percebi que não teríamos a construção automática de DAGs e a funcionalidade de Defer de maneira transparente, essas duas funcionalidades são essenciais para um fluxo de trabalho ágil e eficiente, principalmente tendo em vista que estamos construindo uma plataforma self-service, onde pessoas desenvolvedoras sem experiência com engenharia de dados terão que trabalhar.
No Lakeflow Pipelines (LFP), a estrutura é organizada principalmente em torno de Pipelines. Existem regras e limitações específicas, uma delas é que cada Pipeline suporta no máximo 16 Flows concorrentes, essa é uma limitação que não é documentada.
Em uma DAG tradicional do dbt, encadeamos todos os recursos (tabelas, views e MVs) de forma contínua, sem essa limitação. No LFP, portanto, precisamos planejar melhor quais recursos residem em cada Pipeline e o seu escopo.
Porém, nossa equipe estava habituada a uma experiência de desenvolvimento em que a pessoa não precisa se preocupar em como a DAG será montada. Bastava declarar os relacionamentos entre tabelas via aliases (como ref e source no dbt) e a DAG era gerada de forma transparente no Airflow via Astronomer Cosmos.
Vale abrir um parêntese, embora seja viável utilizar o dbt integrado à Databricks, essa abordagem geraria MVs isoladas (standalone). Consequentemente, cada Materialized View demandaria sua própria pipeline dedicada e a alocação de uma nova instância serverless, elevando consideravelmente o custo computacional. Além disso, a execução da atualização ocorre por meio de um SQL Warehouse, criando uma cobrança duplicada. Para contornar essa proliferação de pipelines, a alternativa seria recorrer à estratégia de cargas incrementais no padrão do dbt, o que nos levaria exatamente de volta ao problema inicial.
Mason - Nosso ‘clone” do DBT para a Databricks
Para resolver essa lacuna, criamos o Mason, nossa própria CLI. Com o auxílio de LLMs na escrita de código, desenvolvemos a primeira versão funcional em apenas duas semanas e, ao longo das duas semanas seguintes, refinamos e adicionamos funcionalidades à ferramenta.
Atualmente, o Mason reproduz com fidelidade o que o dbt faz de melhor, além de agregar melhorias próprias. Não entrarei no detalhe de cada funcionalidade desenvolvida agora, pois trarei esses detalhes em edições futuras da Jornada.
Neste primeiro momento, basta saber que resolvemos tanto a funcionalidade de Defer quanto a construção automática de DAGs para o encadeamento entre Pipelines.
Definição dos Flows e construção de Pipelines
A única configuração necessária é a declaração da fonte (Source). De forma bem parecida com o dbt, utilizamos um arquivo YAML onde especificamos a origem da tabela, permitindo que a CLI construa as referências automaticamente.
Exemplo de arquivo de source do Mason (sources.yml):
sources:
- schema: vendas
tables:
- name: pedidos
description: "Pedidos realizados na plataforma"
identifier: pedidos_2 # nome físico da tabela no Unity Catalog caso você queira um nome diferente
- name: clientes
description: "Cadastro de clientes"
Por padrão, assumimos a leitura a partir do catálogo bronze, mas é possível sobrescrever o catálogo de origem se necessário.
Em seguida, referenciamos a fonte em um dataset de limpeza inicial utilizando o prefixo src:
-- datasets/silver/vendas/pedidos.sql
CREATE OR REFRESH MATERIALIZED VIEW pedidos
REFRESH POLICY INCREMENTAL
CLUSTER BY AUTO
COMMENT 'Pedidos realizados na plataforma, uma linha por pedido, sem os excluídos na origem.'
AS
SELECT
id AS pedido_id,
cliente_id,
CAST(valor_total AS DECIMAL(18, 2)) AS valor_total,
status,
created_at AS pedido_criado_em
FROM ${src_vendas.pedidos}
WHERE deleted_at IS NULL
-- datasets/silver/vendas/clientes.sql
CREATE OR REFRESH MATERIALIZED VIEW clientes
REFRESH POLICY INCREMENTAL
CLUSTER BY AUTO
COMMENT 'Cadastro de clientes, uma linha por cliente.'
AS
SELECT
id AS cliente_id,
nome,
UPPER(uf) AS uf
FROM ${src_vendas.clientes}
Na sequência, construímos as tabelas fato ou dimensão utilizando a sintaxe de ref:
-- datasets/gold/vendas/dim_cliente.sql
CREATE OR REFRESH MATERIALIZED VIEW dim_cliente (
cliente_id BIGINT NOT NULL,
nome STRING,
uf STRING,
regiao STRING,
CONSTRAINT dim_cliente_pk PRIMARY KEY (cliente_id)
)
REFRESH POLICY INCREMENTAL
CLUSTER BY AUTO
COMMENT 'Clientes, uma linha por cliente, com os atributos usados para filtrar e agrupar os pedidos.'
AS
SELECT
cliente_id,
nome,
uf,
CASE
WHEN uf IN ('AC', 'AM', 'AP', 'PA', 'RO', 'RR', 'TO') THEN 'Norte'
WHEN uf IN ('AL', 'BA', 'CE', 'MA', 'PB', 'PE', 'PI', 'RN', 'SE') THEN 'Nordeste'
WHEN uf IN ('DF', 'GO', 'MS', 'MT') THEN 'Centro-Oeste'
WHEN uf IN ('ES', 'MG', 'RJ', 'SP') THEN 'Sudeste'
WHEN uf IN ('PR', 'RS', 'SC') THEN 'Sul'
END AS regiao
FROM ${ref_silver.vendas.clientes}
-- datasets/gold/vendas/fct_pedido.sql
CREATE OR REFRESH MATERIALIZED VIEW fct_pedido (
pedido_id BIGINT NOT NULL,
cliente_id BIGINT,
data_do_pedido DATE,
pedido_criado_em TIMESTAMP,
status STRING,
valor_total DECIMAL(18, 2),
CONSTRAINT fct_pedido_pk PRIMARY KEY (pedido_id)
)
REFRESH POLICY INCREMENTAL
CLUSTER BY AUTO
COMMENT 'Pedidos, uma linha por pedido, com o valor e o cliente que comprou.'
AS
SELECT
pedido_id,
cliente_id,
CAST(pedido_criado_em AS DATE) AS data_do_pedido,
pedido_criado_em,
status,
valor_total
FROM ${ref_silver.vendas.pedidos}
Um detalhe importante na sintaxe do ref é que exigimos a declaração explícita do schema. Essa foi uma mudança intencional em relação ao dbt, enquanto no dbt não é possível ter arquivos de modelos com o mesmo nome em schemas diferentes, com a nossa abordagem conseguimos dar suporte a nomes duplicados em schemas distintos de forma transparente.
Após criar os arquivos SQL dos datasets, a pessoa desenvolvedora precisa apenas rodar o comando:
$ mason scaffold pipelines
resources/vendas__silver.pipeline.yml criado
resources/vendas__silver.schema.yml criado
resources/vendas__gold.pipeline.yml criado
resources/vendas__gold.schema.yml criado
resources/vars.mason.yml criado
5 arquivo(s), 5 alterado(s).
A partir disso, os arquivos de definição das Pipelines são gerados automaticamente:
# GERADO por `mason scaffold pipelines`. Não editar à mão.
resources:
pipelines:
vendas__silver:
name: "${var.run__name_prefix}vendas__silver"
catalog: ${var.dest__vendas__silver__catalog}
schema: ${var.dest__vendas__silver__schema}
serverless: true
continuous: false
# um arquivo por dataset do grupo silver/vendas
libraries:
- glob:
include: ../datasets/silver/vendas/clientes.sql
- glob:
include: ../datasets/silver/vendas/pedidos.sql
# cada ${src_...} vira um parâmetro com o endereço físico do sources.yml
configuration:
src_vendas.clientes: ${var.src__vendas__silver__vendas__clientes}
src_vendas.pedidos: ${var.src__vendas__silver__vendas__pedidos}
# GERADO por `mason scaffold pipelines`. Não editar à mão.
resources:
pipelines:
vendas__gold:
name: "${var.run__name_prefix}vendas__gold"
catalog: ${var.dest__vendas__gold__catalog}
schema: ${var.dest__vendas__gold__schema}
serverless: true
continuous: false
# um arquivo por dataset do grupo gold/vendas
libraries:
- glob:
include: ../datasets/gold/vendas/dim_cliente.sql
- glob:
include: ../datasets/gold/vendas/fct_pedido.sql
# cada ${ref_...} que aponta para outra pipeline vira um parâmetro aqui
configuration:
ref_silver.vendas.clientes: ${var.ref__vendas__gold__silver__vendas__clientes}
ref_silver.vendas.pedidos: ${var.ref__vendas__gold__silver__vendas__pedidos}
Linters
Desenvolvemos também um conjunto de linters para garantir a qualidade, evitar erros comuns de desenvolvimento e dispensar a necessidade de verificação manual por parte das pessoas desenvolvedoras. Em cada artigo, destacarei os linters relevantes para os conceitos apresentados.
Para o fluxo apresentado até aqui, estes são os principais linters ativos:
-
Fonte não declarada: identifica quando uma definição de
srcfaz referência a uma tabela não mapeada nosources.yml(ex:${src_vendas.produtos}). -
Uso correto de src e ref: restringe o uso de
srcapenas à camada silver, impedindo que tabelas da camada gold leiam diretamente da bronze. As permissões de acesso entre camadas viarefsão estritamente validadas. -
Referência sem endereço completo: bloqueia referências simplificadas no estilo do dbt (ex:
${ref_clientes}). É obrigatório especificar os três níveis (camada, schema e nome), o que viabiliza o uso do mesmo nome em schemas distintos. O linter indica os caminhos possíveis em caso de erro. - Referência a dataset inexistente: identifica referências quebradas e sugere nomes similares existentes.
-
Nome de recurso hardcoded: impede referências diretas em SQL (ex:
FROM silver.vendas.pedidos), pois referências manuais não entram na construção da DAG de dependências. - Dependência cíclica: detecta e exibe cadeias de dependência circulares antes da realização do deploy.
-
Estrutura de arquivos: valida se o dataset está localizado no diretório correto:
datasets/<camada>/<schema>/<arquivo>. -
Pipeline desatualizada: recalcula os arquivos de definição e compara com o conteúdo salvo no repositório. Alerta se alguém esquecer de rodar o comando
mason scaffold pipelines. - Limite de 16 flows excedido: indica quando um grupo ultrapassa o limite de 16 flows por Pipeline e fornece o comando para efetuar a divisão em sub-pipelines. Explicarei essa estratégia de divisão em publicações futuras.
As mensagens do linter são formatadas de forma orientativa para facilitar a correção:
$ mason lint
MASON-E002 datasets/silver/vendas/produtos.sql:12
${src_vendas.produtos} não está em sources.yml (schema vendas).
Acrescente a tabela ao grupo do schema que a lê: é de lá que sai o
endereço físico que o resolver emite.
1 achado(s).
Conclusão
A transição para a Databricks aliviou a carga operacional de manter toda a infraestrutura com uma equipe reduzida, enquanto o Lakeflow Pipelines supriu a necessidade de atualizações incrementais que buscávamos no SDP. No entanto, sentimos falta da experiência de desenvolvimento ágil do dbt, essencial para sustentar nosso modelo de plataforma self-service.
O Mason nasceu para preencher essa lacuna, a pessoa desenvolvedora escreve apenas o SQL, declara os arquivos de origem e utiliza os termos de referência (src e ref), deixando toda a orquestração e validação a cargo da CLI e dos linters.
Embora criar e manter uma ferramenta própria exija esforço, com o advento das LLM's ficou fácil de implementar melhorias específicas para o nosso contexto e corrigir bugs sem depender de terceiros.
Nos próximos artigos, falarei mais sobre o recurso de Defer, outras capacidades do Mason e as lições aprendidas ao lidar com cargas incrementais na prática.
Referências
Stack de onde partimos
-
dbt: documentação,
ref,sourcee Defer - Apache Iceberg: site e manutenção e expiração de snapshots
- Trino: site e Trino Gateway
- Starrocks: site
- Cube: site
- Apache Superset: site
- Lakekeeper, catálogo REST do Iceberg: documentação
- Apache Airflow: site e Astronomer Cosmos, que transforma o projeto dbt em DAG do Airflow
- OpenMetadata: site
- Meltano: site
- Kafka Connect: Debezium e Iceberg Sink Connector
Spark Declarative Pipelines e Databricks
- Spark Declarative Pipelines Programming Guide
- Bringing Declarative Pipelines to the Apache Spark Open Source Project, o anúncio da Databricks
- Lakeflow Spark Declarative Pipelines
- Flows
- Materialized views
- Incremental refresh for materialized views
- CREATE MATERIALIZED VIEW e a cláusula REFRESH POLICY
- Limite de 16 tabelas processadas ao mesmo tempo numa pipeline, discussão na comunidade Databricks
- Declarative Automation Bundles resources, onde fica o YAML da pipeline
- Unity Catalog
Top comments (0)