refactor: migra extração BB Ágil de JSON local para Postgres raw - #17
Merged
Conversation
Elimina o checkpoint em arquivo JSON local (dado sensivel sem controle de acesso/backup, e recalculo do fato inteiro a cada run) em favor de tabelas raw + controle no Postgres (bsc_pnab.raw_bbagil_extrato_transacoes, raw_bbagil_subtransacoes, controle_extracao_bbagil_extrato/subtransacoes). Os 13 filtros de negocio que viviam em regras_negocio_bbagil.py (pandas) viram modelos dbt (dbt/minc/models/bsc_pnab_dbt/bronze|silver|gold), seguindo o mesmo padrao ja usado em agentes_dbt. fato_bbagil passa a ser gerado so pelo dbt, rodado automaticamente pelo minc_cosmos_dag ja existente -- a task do Airflow que fazia upsert foi removida. Corrige de quebra o conflito de merge nao resolvido em dbt/minc/models/sources.yml. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
… raw Remove os models bronze/silver/gold de bsc_pnab_dbt e a config correspondente em dbt_project.yml. Prioridade agora e completar a extracao (preencher raw_bbagil_extrato_transacoes/raw_bbagil_subtransacoes com as ~319 mil combinacoes restantes) antes de investir na camada de transformacao -- fica para depois, sobre as mesmas tabelas raw. sources.yml mantido como esta (conflito de merge resolvido continua valendo, e as duas tabelas raw continuam declaradas como source). Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
O commit anterior (23f0df8) so removeu os .sql do bsc_pnab_dbt -- um 'git add -A' com pathspec pra uma pasta ja apagada por 'git rm' anterior falhou silenciosamente pro resto do comando, deixando dbt_project.yml e o docstring da DAG sem reverter (mas ja sem os models .sql, inconsistente). Completa a reversao agora. Tambem remove dags/data_ingest/transferegov_fundo_a_fundo/sql/ (DDL de referencia para raw_gestao_financeira_lancamentos/subtransacoes) -- pasta que tambem foge do padrao do repositorio (a tabela real ja e criada em runtime por ClientPostgresDB, o .sql era so documentacao paralela que nao deveria ter sido commitada). Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
executar_lote rodava todos os itens pendentes (ate ~527 mil, ~18h de execucao) num unico asyncio.gather, e so gravava no Postgres (raw + controle) depois que TUDO terminava. Duas consequencias reais, achadas ao retomar a extracao com os novos planos de acao: zero visibilidade de progresso durante a execucao inteira, e qualquer interrupcao no meio (queda de VPN, restart de container, erro nao tratado) perdia o trabalho inteiro -- pior que o checkpoint em arquivo antigo, que salvava incremental por item. Adiciona tamanho_lote/ao_concluir_lote opcionais em executar_lote: processa em pedacos (default 2000 em extracao_bbagil_dag, via TAMANHO_LOTE_PERSISTENCIA) e persiste cada pedaco assim que termina, via callback sincrono. extracao_beneficiarios_dag continua sem usar essa opcao, comportamento inalterado. Testado (lotes persistidos incrementalmente, e abort em 401/403/429 no meio de um lote preserva os lotes anteriores ja persistidos). Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
get_contas_agencias_programas() sempre buscava agencia/conta de TODOS os planos de acao, ao vivo, mesmo os ja persistidos em execucoes anteriores (raw_planos_acao_dado_bancario era escrita so pra auditoria, nunca lida de volta como cache). Com ~18 mil planos de acao apos os novos programas, isso significava ~18 mil chamadas HTTP repetidas em toda execucao da DAG, so pra descobrir dado que ja se tinha. get_contas_agencias_programas ganha ids_plano_acao_conhecidos (opcional, mantem o modulo Python puro sem acoplar a Postgres): pula a chamada de agencia/conta pra planos ja conhecidos, retorna so os novos. extracao_bbagil_dag._carregar_entes_transferegov le o que ja esta em raw_planos_acao_dado_bancario, passa o conjunto de IDs conhecidos, e funde com os novos encontrados. Testado com mocks (confirma que so chama a API pros planos realmente novos). Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
…odico do BSC Confirmado ao vivo: o BSC aplica um bloqueio temporario (401/403) apos uso sustentado, independente do throttle -- ~20min de chamadas continuas a 10 concorrentes/0.5s, ~31min a 5 concorrentes/1.05s (throttle reduzido mais cedo nessa sessao). Ou seja, e questao de tempo/volume total, nao so concorrencia -- vai se repetir periodicamente durante toda a extracao (~46h estimadas). Com checkpoint em lote (TAMANHO_LOTE_PERSISTENCIA, commit f900134) cada retry e barato -- so refaz o que nao foi persistido desde o ultimo lote -- entao aumenta retries de 3 pra 100 e retry_delay de 5 pra 10 minutos, pra aguentar o padrao de bloqueio/espera a noite toda sem supervisao, em vez de esgotar as tentativas e parar de vez depois de ~15-30min. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
…ios orfa DagRuns concorrentes autenticavam independentemente no SCA com o mesmo client_id/secret, e o SCA invalida o token anterior ao emitir um novo pro mesmo client -- isso derrubava o token de runs paralelas e causava 401/403 sem relacao com o throttle de requisicoes. So faz sentido 1 run ativa. extracao_beneficiarios_dag.py removida: depende do fato_bbagil.parquet que nao existe mais desde a migracao pra Postgres raw, e nenhum model/DAG consome o output das suas 5 tasks -- fica sem produtor e sem consumidor. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
A coluna guardava o id_plano_acao desde sempre (extracao_bbagil_dag.py injeta esse valor na resposta da API), mas o nome "ente" sugeria que era o municipio/estado/orgao real -- na pratica um ente pode ter varios planos de acao (cada um com sua propria conta bancaria segregada), entao o nome confundia a granularidade real da chave. Renomeado via ALTER TABLE nas 4 tabelas do schema bsc_pnab (dados preservados, ja validado rodando a DAG de ponta a ponta: SELECT, JOIN com raw_planos_acao_dado_bancario e INSERT/upsert todos ok com o novo nome). Nao mexi em regras_negocio_bbagil.py (ja orfao, fora de escopo). Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Resolve conflict in dbt/minc/models/sources.yml: mantem raw_bbagil_extrato_transacoes e raw_bbagil_subtransacoes (nova extracao Postgres) lado a lado com fato_bbagil (fonte legada ainda referenciada pelo stg_bbagil desabilitado) e todas as entries LPG/PNAB/territorio trazidas pelo cotas_dbt (Meta 3) do main.
davi-aguiar-vieira
approved these changes
Aug 5, 2026
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Resumo
Migra a extração financeira do BB Gestão Ágil (BSC/SERPRO) de checkpoint em arquivo JSON local para tabelas raw + controle no Postgres, e conclui a extração histórica completa (677 mil combinações plano de ação × período).
Motivação original: os JSONs continham CPF/CNPJ/conta bancária de beneficiário não anonimizados, salvos só no disco de quem roda a DAG (sem backup nem controle de acesso), e o pipeline recalculava o fato inteiro do zero a cada execução varrendo o diretório — apagar um JSON corrompia silenciosamente o resultado final (upsert sobrescreve, não soma).
O que mudou
Arquitetura (extracao_bbagil_dag.py)
extrair_agencias_transferegov→persistir_agencias_contas_transferegov→extrair_extrato_bbagil→extrair_subtransacoes_bbagil. As tasks antigas de consolidação/fato em pandas foram removidas — esta DAG agora só extrai e deposita dado bruto no Postgres.bsc_pnab:raw_bbagil_extrato_transacoes,raw_bbagil_subtransacoes,controle_extracao_bbagil_extrato,controle_extracao_bbagil_subtransacoes(a tabela de controle substitui oPath.exists()como checkpoint — necessária porque ela registra tentativas por plano de ação × período, enquanto a tabela raw só tem linha quando existe transação; sem ela não dá pra distinguir "período vazio confirmado" de "ainda não verificado").execucao_assincrona_bsc.executar_loteganhoutamanho_lote/ao_concluir_loteopcionais: persiste a cada 2.000 itens processados em vez de só no final — evita perder o trabalho inteiro numa interrupção no meio de uma execução de muitas horas.agencias_transferegov.get_contas_agencias_programasganhou cache deids_plano_acao_conhecidos: evita rebuscar agência/conta de planos de ação já persistidos em execuções anteriores (caiu de ~18 mil chamadas HTTP repetidas por execução pra só as novas).max_active_runs=1: DagRuns concorrentes autenticavam independentemente no SCA com o mesmo client_id/secret, e o SCA invalida o token anterior ao emitir um novo — isso derrubava o token de execuções paralelas e causava 401/403 sem relação com throttle de requisições.retries=100/retry_delay=10min: o BSC aplica bloqueio temporário (401/403) após uso sustentado (~20-40min de chamadas contínuas), independente do throttle configurado — é padrão esperado, não falha real, e o checkpoint em lote torna cada retry barato.enterenomeada paraid_plano_acaonas 4 tabelas do schemabsc_pnab: o nome antigo sugeria ser o município/estado real, mas sempre guardou o ID do plano de ação — um mesmo ente pode ter vários planos de ação (cada um com conta bancária segregada própria), então o nome confundia a granularidade real da chave.Removido
extracao_beneficiarios_dag.py: dependia dofato_bbagil.parquetque não existe mais desde esta migração, e nenhum model/DAG consumia o output das suas 5 tasks (CPF list, CNPJ, BPC, CadÚnico, relação trabalhista) — ficava sem produtor e sem consumidor.dags/data_ingest/transferegov_fundo_a_fundo/sql/(DDL de referência que fugia do padrão do repo — as tabelas já são criadas em runtime peloClientPostgresDB).Corrigido de quebra
dbt/minc/models/sources.yml.Resultado da extração (dados reais, schema
bsc_pnab)raw_bbagil_extrato_transacoes(de 677.207 combinações testadas; a diferença retornou "sem lançamentos", registrado como fato de negócio viacontrole_extracao_bbagil_extrato, não erro)raw_bbagil_subtransacoesNão incluído nesta PR (próximos passos)
extracao_beneficiarios_bbagil_dagpro mesmo padrão (precisa antes recriar ofato_bbagilvia dbt)Test plan
dag_run.state=successraw_planos_acao_dado_bancarioe INSERT/upsert confirmados funcionando comid_plano_acaomax_active_runs=1confirmado no Postgres (SELECT max_active_runs FROM dag) e sem mais 401/403 por concorrência entre DagRuns