Skip to content

[bug] ajuste no sistema de atualização #1781

Description

@laura-l-amaral

Descrição do bug

  • verificar se a max_date é a cobertura ou a data de atualização
  • incluir parametro de cobertura ou data_atualização na poll_source_for_update_task
  • comparar bananas com bananas e laranjas com laranjas

Como reproduzir

https://basedosdados.org/dataset/562b56a3-0b01-4735-a049-eeac5681f056?table=95106d6f-e36e-4fed-b8e9-99c41cd99ecf

Referências

No response


Investigação

O poll atual (poll_source_for_update, em pipelines/utils/metadata/register.py) sempre comparava a data da fonte contra Table.Update.latest (via client.get_table_update_latest()) — um campo que pode ser um timestamp de execução (bq.last_modified), não a cobertura real dos dados. Foi exatamente esse defeito, já confirmado em produção, que causou os bugs de detecção de br_me_caged (#1760) e br_ans_beneficiario (#1779).

Criamos um worktree do repositório no commit imediatamente anterior à migração do CAGED para Prefect 3 (1dff7cfc) para inspecionar o sistema de poll do Prefect 0. A função antiga, check_if_data_is_outdated, já tinha exatamente o parâmetro que esta issue pede:

def check_if_data_is_outdated(
    dataset_id, table_id, data_source_max_date,
    date_type: str = "data_max_date",
    date_format: str = "%Y-%m-%d",
) -> bool:
    if date_type == "data_max_date":
        data_api = get_api_most_recent_date(...)   # Coverage.DateTimeRange — cobertura real
    if date_type == "last_update_date":
        data_api = get_api_last_update_date(...)   # allUpdate(table_Id).latest — Table.Update.latest
    ...

Auditoria completa (grep de todas as chamadas a check_if_data_is_outdated no worktree pré-migração, com parsing de parênteses balanceados pra extrair o date_type de cada uma) mostrou que 19 dos ~23 datasets auditados usavam o default (data_max_date, correto). Esse parâmetro se perdeu na migração para Prefect 3 — poll_source_for_update_task não tinha escolha nenhuma, sempre comparava contra Table.Update.latest. Não foi uma decisão consciente, foi uma regressão silenciosa da migração.

Correção — mergeado em #1783

  • Restaura a escolha: poll_source_for_update(..., compare_against="coverage") (novo parâmetro; nomes "coverage"/"table_update", diferentes do antigo date_type por serem menos ambíguos). Novo método MetadataClient.get_coverage_max_date() lê Coverage.DateTimeRange (mesma fonte que o antigo get_api_most_recent_date consultava). Default é "coverage" — o mais usado entre os flows auditados.
  • Aplicado explicitamente nos 32 pontos de chamada de poll_source_for_update_task que existem hoje no repositório — não foi só uma mudança na função compartilhada: em cada um dos ~30 arquivos de flow, adicionamos o argumento compare_against na chamada, decidido caso a caso (histórico do Prefect 0 quando existia, natureza da fonte quando não existia — 9 datasets sem esse histórico, ver auditoria completa no PR).
  • commit_source_update_task (grava RawDataSource.Update.latest) movido de posição em todos esses mesmos ~31 arquivos: era chamado só no fim do flow, depois de register_table_materialization_task — ou seja, depois de baixar, subir e materializar tudo. Passou a ser chamado logo depois do poll confirmar que há dado novo, antes de baixar/materializar. Assim, se o flow falha no meio (download, upload, dbt), o RawDataSource.Update já reflete que a fonte publicou — antes ficava sem nenhuma pista de que havia novidade, só o registro de que a run falhou. A função também ganhou dois parâmetros novos, update_metadata e materialize_after_dump (default True), e passou a decidir sozinha se grava — cada flow só passa os dois flags que já tem, sem precisar embrulhar a chamada num if próprio (evita repetir esse guard nos ~25 flows que chamam a task). register_table_materialization_task continua no fim, sem mudança.

Achado extra: cobertura NonHistorical não tem baseline pra comparar

Um comentário de code review em br_me_cnpj/flows.py apontou que a tabela simples usa cobertura NonHistorical() — usada em tabelas sem coluna de data confiável, onde a única cobertura possível vem do metadado last_modified do próprio BigQuery. register_table_materialization, para esse tier, só grava Table.Update.latest — nunca grava Coverage.DateTimeRange. Com compare_against="coverage" (o default novo), get_coverage_max_date sempre retornaria None pra essas tabelas, e a lógica trata "sem baseline" como "sempre há atualização" — ou seja, o poll nunca diria "sem novidade", rodando o flow todo dia mesmo sem nada novo na fonte.

Conferimos contra o Prefect 0: lá, simples também não passava nenhum date_type (mesmo default, data_max_date/coverage), mas funcionava porque a rotina antiga de materialização, com historical_database=False, gravava uma Coverage.DateTimeRange "falsa" (tipo all_free, start=end) derivada do mesmo last_modified do BigQuery — ou seja, comparava exatamente o mesmo valor, só que guardado num campo diferente. No Prefect 3, esse valor já está em Table.Update.latest. Corrigido usando compare_against="table_update" só nas tabelas NonHistorical (br_me_cnpj/simples; br_rf_cnpj/simples e dicionario, que tem o mesmo padrão), mantendo "coverage" nas demais tabelas de cada dataset.

Estado final de cada flow

compare_against="coverage" (27 pontos de chamada — competência real na fonte)

br_anatel_banda_larga_fixa, br_anatel_telefonia_movel, br_cgu_cartao_pagamento, br_cgu_servidores_executivo_federal, br_cgu_licitacao_contrato, br_cgu_beneficios_cidadao, br_ms_sia, br_ms_sih, br_ms_sinan, br_ibge_inpc, br_ibge_ipca, br_ibge_ipca15, br_me_cnpj (empresas/socios/estabelecimentos), br_rf_cnpj (empresas/socios/estabelecimentos), br_rf_cno, br_anp_precos_combustiveis, br_bcb_agencia, br_bcb_estban, br_denatran_frota, br_inmet_bdmep, br_rf_cafir, br_stf_corte_aberta, br_me_caged (microdados_movimentacao/_fora_prazo/_excluida — ver nota abaixo).

Sem histórico no Prefect 0 (datasets novos, decisão pela natureza da fonte): br_bcb_taxa_selic, br_me_siconfi, br_mf_divida_ativa, au_abs_cpi, au_abs_labour_force, us_bls_cpi, us_bls_qcew, world_cricsheet.

compare_against="table_update" (8 pontos de chamada — exceções deliberadas, não bugs)

br_cvm_fi, br_tse_eleicoes, br_cgu_emendas_parlamentares, br_ibge_pnadc (já usavam last_update_date de propósito no Prefect 0), br_bndes_operacoes_contratadas (2 flows, sem histórico Prefect 0 — source_max_date é o last_modified de um recurso CKAN, timestamp de publicação, enquanto a Coverage da tabela é anual; comparar contra Coverage misturaria granularidades desconexas), e br_me_cnpj/simples, br_rf_cnpj/simples,dicionario (cobertura NonHistorical — ver achado extra acima).

br_me_caged — nota especial, PR #1760 fechado

O #1760 propunha migrar br_me_caged para um modelo de poll totalmente novo (register_source_coverage_task/check_source_is_ahead_of_table_task/sync_table_coverage_task, comparando RawDataSource.Update contra Table.Update repropositado pra guardar cobertura em vez de timestamp de execução). Em vez disso, aplicamos compare_against="coverage" direto no br_me_caged (igual aos outros 26 flows acima) — resolve o mesmo problema, mais simples, sem repropositar Table.Update.

Durante o teste final do #1783, achamos evidência ao vivo de que a abordagem do #1760 é mais frágil: o deployment de br_me_caged__microdados_movimentacao_excluida (rodando a branch do #1760) tinha Table.Update.latest contaminado com um timestamp de execução (2026-08-12, não uma data de cobertura), bloqueando check_source_is_ahead_of_table_task de detectar novidade — porque Table.Update é um campo compartilhado, e register_table_materialization_task (usado em todos os outros ~30 flows) ainda escreve nele com a semântica antiga, contaminando de volta. Coverage.DateTimeRange não tem esse ponto de quebra. #1760 foi fechado em favor de compare_against="coverage", já testado ao vivo em produção nas 3 tabelas do CAGED (sem duplicação de dados — verificado via BigQuery depois de um teste com execução concorrente acidental).

Pendência

Datasets com compare_against="coverage" podem ainda ter Table.Update.latest contaminado com um timestamp de execução "adiantado" (mesmo cenário do CAGED/ANS antes da correção manual de metadados) — o #1783 só trocou o alvo da comparação, não corrigiu retroativamente valores já contaminados. Vale uma checagem individual em cada um.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

bugDefeito em código, tooling, CI ou comportamento de pipeline

Type

No type

Projects

Milestone

No milestone

Relationships

None yet

Development

No branches or pull requests

Issue actions