o projeto

1.01  projeto

Pipeline de dados financeiros da CVM

● em curso  ·  bronze e silver agendadas, gold em construção

Como entregar demonstrações financeiras prontas para análise, se a fonte mistura escalas, republica exercícios sem dizer o que mudou e organiza as contas numa hierarquia que o consumidor teria de remontar a cada consulta?

os dados

O que entra no pipeline

Toda companhia aberta entrega à CVM, a Comissão de Valores Mobiliários, as demonstrações financeiras do exercício anterior, uma vez por ano. A CVM reúne essas entregas em pacotes anuais e os disponibiliza para download direto, sem cadastro. Hoje o pipeline processa três desses pacotes.

3demonstrações
~700companhias por ano
6exercícios, 2021 a 2026
~2,6 miregistros na bronze
~1 miregistros na silver

DRE · demonstração do resultado do exercício

Receitas, custos, despesas e lucro. É uma demonstração de fluxo: tem data de início e de fim do exercício, e responde quanto a companhia ganhou ou perdeu no período.

Bronze 101_dre_dfp, 428.532 registros.
Silver 201_dre_dfp, 163.013 registros, 24 colunas.

BPA · balanço patrimonial ativo

Bens e direitos da companhia num ponto do tempo. Por ser uma fotografia da posição, e não um fluxo, não existe data de início. O pipeline trata essa diferença: a tabela da BPA não tem a coluna de início do exercício que a DRE tem.

Bronze 102_bpa_dfp, 815.877 registros.
Silver 202_bpa_dfp, 309.246 registros. DDL sem DT_INI_EXERC.

BPP · balanço patrimonial passivo

Obrigações e patrimônio líquido. Junto com o ativo, forma o balanço completo, e segue a mesma estrutura e a mesma lógica de processamento da BPA.

Bronze 103_bpp_dfp, 1.394.995 registros.
Silver 203_bpp_dfp, 528.874 registros.

A CVM publica outras três demonstrações: DFC, DVA e DMPL. A estrutura modular as acomoda sem mudança de arquitetura, seguindo a numeração 104/204, 105/205 e 106/206.

termosdemonstrações financeirasfluxo e posiçãoDDL

arquitetura

Arquitetura medalhão

O dado passa por estágios com papéis definidos. O que entra bruto no primeiro sai tipado, normalizado e enriquecido na silver, pronto para análise sem que o consumidor precise saber de onde veio ou como foi limpo.

○ ainda não construída

landing zone · os arquivos como a CVM publicou

Os pacotes originais ficam guardados exatamente como chegaram, cada um acompanhado de um registro de quando a CVM o alterou pela última vez. Antes de qualquer substituição, a versão anterior é arquivada com carimbo de tempo.

Uma rotina diária pergunta à CVM, exercício por exercício, se o pacote mudou, e só baixa o que mudou. Os notebooks de processamento nunca fazem chamada de rede à origem.

ZIPs em Unity Catalog Volumes, com _metadata.json guardando o cabeçalho Last-Modified da origem. Verificação por requisição HEAD; versões anteriores em archive/.

/Volumes/.../landing/dfp/
dfp/
├── 2023/
│   ├── dfp_cia_aberta_2023.zip
│   ├── _metadata.json
│   └── archive/
│       └── dfp_cia_aberta_2023_20260815_063012.zip
├── 2024/
│   ├── dfp_cia_aberta_2024.zip
│   └── _metadata.json

bronze · dado bruto com rastreio completo

Preserva a estrutura original da CVM, sem nenhuma transformação. Cada ingestão acrescenta quatro colunas de rastreio e soma linhas à tabela, sem apagar nada: nenhuma versão publicada pela CVM se perde.

Processar duas vezes a mesma versão não gera duplicata, porque antes de começar o notebook confere se aquela combinação de fonte, ano, data de alteração e status já foi registrada.

Gravação append-only. Idempotência garantida pela tabela controle_ingestao.

01_bronze/101_cvm_dfp_dre.py
df_bronze = df_validado \
    .withColumn("_versao_ingestao", lit(versao_atual)) \
    .withColumn("_last_modified_cvm", lit(last_modified_cvm)) \
    .withColumn("_ingest_ts", current_timestamp()) \
    .withColumn("_source_file", lit(f"dfp_cia_aberta_DRE_con_{ano}.csv"))

silver · dado padronizado e enriquecido

Seleciona a versão mais recente de cada registro na bronze, converte os tipos, normaliza a escala monetária e deriva cinco colunas a partir da estrutura contábil.

A gravação substitui o ano inteiro numa única operação: ou tudo muda, ou nada muda. Isso eliminou o instante de tabela incompleta que existia quando a escrita era feita apagando e reinserindo.

Versão vigente resolvida por window function. Gravação atômica com replaceWhere, no lugar do antigo DELETE seguido de APPEND.

02_silver/201_cvm_dfp_dre.py
df_silver.write \
    .format("delta") \
    .mode("overwrite") \
    .option("replaceWhere", f"ANO = {ano}") \
    .saveAsTable(f"{SCHEMA_SILVER}.201_dre_dfp")

gold · análise de negócio, em construção

Métricas de negócio, indicadores por companhia, comparações setoriais e evolução ao longo do tempo. É o primeiro item da fila, e só fazia sentido depois da silver normalizada, que ficou pronta em setembro.

16tabelas Delta
15notebooks
24colunas na silver
3jobs agendados
termosarquitetura medalhãoappend-onlyidempotênciaescrita atômica

negócio

Três problemas que a silver resolve

Os dados saem da CVM com problemas que a fonte não corrige. O pipeline os resolve na silver, antes que cheguem a qualquer consumidor.

Escala monetária misturada

o problema Cerca de 540 companhias reportam em milhares e 15 em unidades. O valor chega em um campo e a escala em outro, e a coluna de valor não distingue as duas. Sem conversão, qualquer soma ou ranking entre elas está errado por um fator de mil.

arraste para aplicar a escala simulação
Companhia A · lucro de 2024
× 1.000, escala aplicada
1.234.567.000

Companhia B · lucro de 2024
já em unidades
890.412.300 com a escala aplicada, A é maior que B.
a ordem entre as duas se inverte.
Companhia A · lucro de 2024
escala MIL
1.234.567

Companhia B · lucro de 2024
escala UNIDADE
890.412.300 lado a lado, B parece setecentas
vezes maior do que A.
como chegacom a escala aplicada

Republicação silenciosa

o problema Uma companhia entrega a DRE de 2024 em março. Em julho, reentrega com o lucro corrigido. A CVM substitui o pacote do exercício sem dizer o que mudou. O pipeline guarda todas as versões na bronze e resolve a vigente na silver.

lucro do exercício de 2024 simulação
mar 2025dez 2025
consultando em dez 2025, o pipeline responderia 1.234.567 versão 3, entregue em 18 nov 2025
versãoentregue emlucro divulgado
112 mar 20251.198.402
204 jul 20251.221.870
318 nov 20251.234.567

como o versionamento funciona na prática

Cada versão entra na bronze com um número de ingestão crescente. A silver agrupa os registros pela chave de negócio e fica só com o de número mais alto. O resultado é que a silver sempre tem o valor vigente, enquanto a bronze guarda o histórico completo: dá para comparar a versão de março com a de julho.

02_silver/201_cvm_dfp_dre.py
window_spec = Window.partitionBy(
    "CNPJ_CIA", "DT_REFER", "CD_CONTA", "ORDEM_EXERC"
).orderBy(col("_versao_ingestao").desc())

df_versao_atual = df_bronze \
    .withColumn("_row_num", row_number().over(window_spec)) \
    .filter(col("_row_num") == 1)

Estrutura hierárquica contábil

o problema A CVM publica as contas numa hierarquia escrita com pontos: 3, depois 3.01, depois 3.01.01. O valor do pai é a soma dos filhos. Sem enriquecimento, quem consulta precisa remontar essa estrutura toda vez. A silver deriva cinco colunas que tornam a hierarquia consultável diretamente.

exemplo de hierarquia · DRE estrutura real
  • 3Resultado Bruto nível 1 · totalizadora · raiz
    • 3.01Receita Líquida nível 2 · totalizadora · aditiva
      • 3.01.01Receita de Vendanível 3 · analítica
      • 3.01.02Deduçõesnível 3 · analítica
    • 3.02Custo dos Bens Vendidos nível 2 · totalizadora · derivada
CD_CONTADS_CONTANIVELPAIRAIZTIPOESTRUT.
3Resultado Bruto1NULL3TOTALIZ.ADITIVA
3.01Receita Líquida233TOTALIZ.ADITIVA
3.01.01Receita de Venda33.013ANALÍT.NULL
3.01.02Deduções33.013ANALÍT.NULL
3.02Custo dos Bens Vend.233TOTALIZ.DERIVADA

as cinco colunas derivadas na silver

NIVEL_CONTA
Profundidade na hierarquia, de 1 a 5, calculada pela contagem de pontos no código.
CD_CONTA_PAI
Tudo o que vem antes do último ponto: 3.01.01 aponta para 3.01. Nulo no nível 1.
CD_CONTA_RAIZ
O primeiro segmento do código: 3.01.01 pertence à raiz 3.
TIPO_CONTA
Totalizadora, até dois níveis, ou analítica, a partir de três. Um painel precisa filtrar por nível para não somar a mesma coisa duas vezes.
TIPO_ESTRUTURAL
Aditiva, quando o pai é a soma dos filhos, ou derivada, quando a conta é calculada, como 3.02 e 3.03. A diferença importa para as validações contábeis.
termosnormalizaçãowindow functionversionamentohierarquia contábil

qualidade

Guardrails e observabilidade

As verificações acontecem antes de qualquer tabela ser modificada. Se uma delas falha, os dados que já estavam lá ficam intactos, e a falha fica registrada para consulta.

o que acontece quando um arquivo chega com problema simulação

quatro tabelas de observabilidade

Cada execução deixa registro em tabela, e não em log que rola para fora da tela. Assim, perguntas como quantos registros vieram naquela noite ou em que ponto a rotina parou se respondem com uma consulta.

controle de ingestão
Cada ingestão por fonte e ano, com a data de alteração na CVM e o status. É a base da detecção de anos pendentes.
execuções
Cada execução de notebook, com job, run, task, duração e contagem de registros.
guardrails
Cada verificação, com o tipo, o resultado (aprovada, falhou ou alerta) e os detalhes.
jobs
Cada run de job, com início gravado pela primeira task e fim pela última. O status só piora: uma vez erro, sempre erro.

controle_ingestao, observabilidade_execucoes, observabilidade_guardrails e observabilidade_jobs, esta última atualizada por MERGE idempotente.

reconciliação em quatro pontos

A bronze confere a contagem de registros em quatro momentos. Se o número muda entre dois deles, a divergência é registrada e o processamento daquele ano para, sem afetar os demais.

01_bronze/101_cvm_dfp_dre.py
count_extraido = len(df_pandas)           # 1. extração do CSV
count_spark    = df_raw.count()            # 2. conversão para Spark
count_validado = df_validado.count()       # 3. pós-validação de schema
count_tabela   = spark.table(...).count()  # 4. pós-gravação na Delta
termosguardrails de qualidadefail-fastreconciliaçãoobservabilidade

operação

Como o pipeline roda

Três jobs separados, cada um com um papel, declarados em arquivo versionado. Perguntar à CVM se algo mudou é tarefa diária e leve; processar é tarefa pesada e semanal; testar acontece quando o código muda.

Verificação diária

todo dia, às 6h

Pergunta à CVM, exercício por exercício, se o pacote mudou desde a última vez. Quando mudou, arquiva a versão atual e baixa a nova. Nada do que já estava ali é sobrescrito sem cópia.

Notebook 004_verificacao_diaria. Requisição HEAD e comparação do Last-Modified com os metadados locais.

termosdetecção de mudançaagendamento

como o pipeline decide o que processar

Antes de processar, o pipeline consolida o que está pendente nas seis tabelas de destino, três da bronze e três da silver, cruzando com a tabela de controle. O resultado é a lista de anos a processar, limitada aos exercícios mais recentes.

Quem roda pode sobrescrever essa decisão, nesta ordem de prioridade: parâmetro do notebook, variável de ambiente, detecção automática e, por último, o ano atual. Uma carga completa força todos os anos disponíveis.

Lista ANOS_PROCESSAR; carga completa com CARGA=completa.

infraestrutura

Infraestrutura como código

O ambiente e os jobs estão declarados no repositório. Quem clonar sobe o mesmo pipeline sem depender de nada configurado na interface. Os dados não vêm junto: o pipeline os baixa da CVM na execução.

três ambientes
Desenvolvimento, teste e produção declarados num único arquivo, que compõe os nomes dos schemas. Fora do Databricks, a configuração assume desenvolvimento, para os testes rodarem no CI sem quebrar.
asset bundles
Jobs, tasks, agendamentos e parâmetros em YAML versionado, publicados com um comando por ambiente.
integração contínua
Dois fluxos no GitHub Actions: padrão de código e testes unitários a cada push, e o teste ponta a ponta no Databricks, sob demanda.
serverless
Leitura e escrita direto nos volumes do Unity Catalog, sem área temporária e sem depender de utilitários do cluster. Funciona em Spark Connect.

como os schemas são compostos

Produção está declarada, mas não instanciada. A configuração recusa esse ambiente com uma mensagem explícita. Sem consumidor real, manter um ambiente de produção separado seria custo sem retorno, e a arquitetura o comporta sem mudança de código quando fizer sentido.

ambientes.json
{
  "composicao": { "schema": "{prefixo}_{ordem}_{camada}" },
  "ambientes": {
    "dev":  { "prefixo": "proj_cvm_dev",  "existe": true  },
    "test": { "prefixo": "proj_cvm_test", "existe": true  },
    "prod": { "prefixo": "proj_cvm_prod", "existe": false }
  }
}
// dev  → workspace.proj_cvm_dev_01_bronze.101_dre_dfp
// test → workspace.proj_cvm_test_01_bronze.101_dre_dfp
stackDatabricksPySparkDelta LakeUnity CatalogPythonGitHub ActionsRuffpytest

registro

Decisões

Cada decisão fica no registro de evolução do repositório, inclusive as que desfizeram uma escolha anterior, marcadas em terracota.

  • 24 ago

    Passei a testar o pipeline no Databricks de verdade

    Testes de integração ponta a ponta: dez tasks, schemas isolados, cinco validações por demonstração.

  • 23 ago

    Coloquei verificação automática em cada alteração

    GitHub Actions com Ruff para padrão de código e pytest para as funções puras.

  • 19 ago

    Parei de configurar o ambiente na mão

    Migração para Databricks Asset Bundles, com jobs declarados em YAML e três targets.

  • 19 ago

    Troquei a forma de regravar a silver

    Gravação atômica com replaceWhere, eliminando a janela de tabela incompleta do DELETE seguido de APPEND.

  • 17 ago

    Descobri que uma verificação minha passava sempre

    Premissa oculta em validação hierárquica: a amostra não era representativa, e nela a checagem não tinha como falhar.

  • 10 ago

    Acrescentei a terceira demonstração

    BPP adicionada no padrão modular, 103 e 203, com a mesma arquitetura das anteriores.

  • 03 ago

    Tornei a bronze idempotente

    Gravação append-only com versionamento, eliminando as duplicatas estruturais.

  • 31 jul

    Adaptei o pipeline para rodar em serverless

    Refatoração para Spark Connect, sem área temporária em /tmp.

  • 31 jul

    Documentei uma limitação real do Delta Lake

    ALTER COLUMN TYPE não é suportado, o que muda a forma de evoluir o schema.

  • 22 jul

    Criei a landing zone com detecção de mudança

    Arquivos em Unity Catalog Volumes, com detecção de atualização pelo cabeçalho Last-Modified.

próximos passos

O que está aberto

em construção
Camada gold: métricas de negócio, indicadores por companhia e rankings setoriais. É o primeiro item da fila.
na fila
Visualização dos dados da gold, em Genie Spaces ou em qualquer ferramenta de BI sobre o Unity Catalog.
na fila
Mais três demonstrações: fluxo de caixa (DFC), valor adicionado (DVA) e mutações do patrimônio líquido (DMPL).

repositório