Conversation
… vs timestamp de execução) O poll antigo comparava a data publicada pela fonte (FTP) contra Table.Update.latest, que é o timestamp de quando a tabela foi materializada (bq.last_modified), não a competência coberta pelos dados. Isso trava a detecção de atualização sempre que uma materialização anterior grava um timestamp de execução "adiantado" em relação ao dia-1 do mês publicado pela fonte seguinte — foi o que aconteceu com microdados_movimentacao: Table.Update ficou em 2026-07-07 e bloqueou a detecção de junho/2026, já disponível no FTP desde 29/07. Migra os 3 flows do dataset (compartilham _run_me_caged) para o modelo de poll novo (pipelines/utils/metadata/poll.py), já usado pelo CNES: register_source_coverage_task + check_source_is_ahead_of_table_task antes do download, sync_table_coverage_task depois da materialização — que grava Table.Update com a cobertura real lida do BigQuery, não com o horário de execução.
📝 WalkthroughWalkthroughThe CAGED flow registers source coverage, checks source freshness against table coverage, and synchronizes table coverage. It removes direct source polling, table materialization registration, and the conditional source update commit. ChangesCAGED coverage flow
Estimated code review effort: 2 (Simple) | ~10 minutes Possibly related issues
Possibly related PRs
Suggested labels: Suggested reviewers: 🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@pipelines/datasets/br_me_caged/flows.py`:
- Around line 48-60: Before enabling the scheduled flow, repair the production
Table.Update.latest metadata that is stale or contaminated: either add the
required correction to this flow before the source-ahead gate, or perform a
manual recovery run with force_run=True and update_metadata=True. Ensure the
correction occurs before normal scheduled runs can return early, while
preserving the existing materialization and coverage behavior.
🪄 Autofix
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: defaults
Review profile: CHILL
Plan: Pro Plus
Run ID: 1fb89040-9f0d-4ef1-b71a-d04f783680df
📒 Files selected for processing (1)
pipelines/datasets/br_me_caged/flows.py
Marca as 3 chamadas de metadata (register_source_coverage_task, check_source_is_ahead_of_table_task, sync_table_coverage_task) que hoje apontam direto para env="prod", para facilitar testar a migração do poll novo contra o backend de dev antes de tocar nos registros de produção.
|
Note GitHub couldn't provide a complete incremental comparison for this pull request, so CodeRabbit is performing a full review instead. This review may take a little longer. |
Validação em produção (2026-08-09)Rodei os 3 flows do dataset a partir desta branch para validar a correção antes do merge:
Conferido direto no BigQuery (fora do que o próprio flow logou) — junho/2026 chegou certo nas 3 tabelas, sem duplicação:
Pendência: essa validação rodou a partir desta branch, não do deploy normal via |
… pro poll.py) Teste da abordagem alternativa ao PR #1760 (migração completa pro poll.py): br_me_caged ainda usa poll_source_for_update_task, mas agora com compare_against="coverage" em vez do default antigo (Table.Update.latest, o campo contaminado que causou o bug original). register_table_materialization_task já atualiza Coverage.DateTimeRange a partir do BigQuery corretamente — isso nunca foi o problema. Só a comparação de leitura estava errada. Se compare_against="coverage" sozinho resolver, não precisamos da migração inteira pro poll.py. Afeta os 3 flows (microdados_movimentacao, _fora_prazo, _excluida) via _run_me_caged compartilhado. Teste-piloto: microdados_movimentacao_fora_prazo.
Fecha a causa raiz sistêmica por trás da issue #1781: poll_source_for_update sempre comparava a fonte contra Table.Update.latest (timestamp de execução), nunca contra Coverage.DateTimeRange (cobertura real) — um parâmetro que existia no Prefect 0 (date_type) e se perdeu silenciosamente na migração pro Prefect 3. - Restaura a escolha como compare_against="coverage"|"table_update" em poll_source_for_update/poll_source_for_update_task, com validação (ValueError em valor inválido) e logging das datas comparadas. - MetadataClient.get_coverage_max_date(): novo método de leitura de Coverage.DateTimeRange. - Aplicado explicitamente nos 32 flows que usam poll_source_for_update_task, auditados um a um contra o comportamento do Prefect 0 (worktree pré-migração); default trocado de "table_update" para "coverage" depois que todos os callers passaram a ser explícitos. - Corrige br_me_cnpj/simples e br_rf_cnpj/simples,dicionario: tabelas NonHistorical não têm Coverage.DateTimeRange, então precisam de compare_against="table_update" especificamente. - commit_source_update_task movido para logo após o poll confirmar dado novo (antes de baixar/materializar), em vez de só no fim do flow — se o flow falhar no meio, o RawDataSource.Update ainda reflete que a fonte publicou. Ganhou os parâmetros update_metadata/materialize_after_dump (default True) para decidir sozinha se grava, evitando repetir o mesmo guard em ~25 flows. register_table_materialization_task continua no fim, sem mudança. - Corrige o Pyrefly type check (quebrado por um PR não relacionado, onboarding do us_harvard_cbdb) adicionando o exclude correspondente. - Testado ao vivo em produção nas 3 tabelas do br_me_caged (compare_against e o commit adiantado, confirmados via logs e BigQuery). Relacionado: #1784 (consolidado nesta branch), #1760 (decisão em aberto, compare_against="coverage" pode tornar a migração pro poll.py innecessária).
|
Fechando este PR — o mesmo problema (poll comparando contra Durante os testes finais do #1783 usando exatamente as tabelas do CAGED, achamos evidência ao vivo de que a abordagem deste PR é mais frágil:
|
Contexto
br_me_caged__microdados_movimentacaoestava reportando "Não há novas atualizações na fonte original" mesmo com a fonte (FTPftp.mtps.gov.br/pdet/microdados/NOVO CAGED) já tendo publicado a competência de junho/2026 desde 29/07.Causa raiz
O poll antigo (
poll_source_for_update) compara a data da fonte contraTable.Update.latest— que não é a cobertura real dos dados, e sim o timestamp de quando a tabela foi materializada (bq.last_modified). Uma materialização anterior gravouTable.Update.latest = 2026-07-07, um timestamp de execução. Como a data da fonte é sempre representada pelo dia 1 do mês (2026-06-01), a comparação2026-06-01 > 2026-07-07é falsa — bloqueando a detecção mesmo com dado novo disponível. É o mesmo defeito arquitetural já corrigido nobr_ms_cnes, migrando para o modelo de poll novo (pipelines/utils/metadata/poll.py).O que muda
Os 3 flows do dataset (
microdados_movimentacao,_fora_prazo,_excluida— compartilham_run_me_caged) passam a usar:register_source_coverage_task+check_source_is_ahead_of_table_taskantes do download (em vez depoll_source_for_update_task)sync_table_coverage_taskdepois da materialização (em vez deregister_table_materialization_task+commit_source_update_task)sync_table_coverage_taskgravaTable.Update.latestcom a cobertura real lida do BigQuery (bq.read_max_date), não com o horário de execução — eliminando a causa do travamento.Observação:
Table.Update.latestda tabelamicrodados_movimentacaoestá hoje contaminado com2026-07-07(resíduo do poll antigo) e não vai se autocorrigir sem uma materialização bem-sucedida — precisa de correção manual ou de umforce_run=Trueantes que o gate novo funcione corretamente.