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.
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.
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/.
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.
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.
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.
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.
× 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.
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.
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.
| versão | entregue em | lucro divulgado |
|---|---|---|
| 1 | 12 mar 2025 | 1.198.402 |
| 2 | 04 jul 2025 | 1.221.870 |
| 3 | 18 nov 2025 | 1.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.
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.
- 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
- 3.01Receita Líquida
nível 2 · totalizadora · aditiva
| CD_CONTA | DS_CONTA | NIVEL | PAI | RAIZ | TIPO | ESTRUT. |
|---|---|---|---|---|---|---|
| 3 | Resultado Bruto | 1 | NULL | 3 | TOTALIZ. | ADITIVA |
| 3.01 | Receita Líquida | 2 | 3 | 3 | TOTALIZ. | ADITIVA |
| 3.01.01 | Receita de Venda | 3 | 3.01 | 3 | ANALÍT. | NULL |
| 3.01.02 | Deduções | 3 | 3.01 | 3 | ANALÍT. | NULL |
| 3.02 | Custo dos Bens Vend. | 2 | 3 | 3 | TOTALIZ. | DERIVADA |
as cinco colunas derivadas na silver
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.
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_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.
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
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.
Pipeline semanal
toda segunda-feira, às 7h
Processa bronze e silver das três demonstrações em trilhos paralelos. DRE, BPA e BPP rodam ao mesmo tempo, e cada par bronze e silver é uma cadeia independente: se uma demonstração falha, as outras seguem.
Sete tasks, execução paralela, timeout de duas horas.
Testes de integração
sob demanda
Sobe o pipeline completo num ambiente isolado e o executa sobre o exercício de 2021. Para cada demonstração, confere cinco coisas: se a tabela existe, a contagem, a reconciliação, a unicidade das chaves e os metadados.
Dez tasks em schemas proj_cvm_test_*, disparadas pelo GitHub Actions, com interrupção na primeira falha.
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.
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.
{
"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
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