Class BulkUpserter

java.lang.Object
br.com.vantroba.bi.transporte.processamento.BulkUpserter

@Singleton public class BulkUpserter extends Object
Caminho idempotente UPSERT diff-aware + DELETE seletivo (spec 004).

É o caminho de gravação canônico do orquestrador (substitui o legado DELETE+COPY removido na fase 6): o COPY vai para uma TEMP TABLE (ON COMMIT DROP, fora da publicação de replicação), o INSERT … ON CONFLICT … DO UPDATE … WHERE IS DISTINCT FROM … preserva id e xmin quando o conteúdo já bate, e o DELETE … WHERE … NOT EXISTS … apaga apenas o que sumiu do upload (não o período inteiro).

Métricas

O RETURNING (xmax = 0) AS inserted distingue INSERT puro (true) de UPDATE in-place (false). Linhas filtradas pelo WHERE IS DISTINCT FROM não aparecem no RETURNING — viram UpsertResultado.inalteradas().

Empty-guard

Lista vazia retorna UpsertResultado(0,0,0,0) sem rodar UPSERT nem DELETE seletivo — caso contrário, NOT EXISTS (… tabela vazia …) apagaria o período inteiro (spec § 4.6).

Preservação de campos de cadastro

As colunas de ColunasConteudo.PRESERVAR_SE_VAZIO_NA_ENTRADA (cadastro de veículo/frota) recebem tratamento especial: quando o valor de entrada é "vazio" (NULL, string vazia ou 0), o UPSERT mantém o valor já gravado em vez de sobrescrevê-lo. Tanto o SET quanto o WHERE IS DISTINCT FROM usam COALESCE(NULLIF(EXCLUDED.col, <sentinela>), <tabela>.col) — assim uma entrada vazia não dispara UPDATE (sem churn de xmin/replicação) e correções de cadastro feitas direto no banco sobrevivem a reprocessamentos de Excel. Ver valorEfetivo(String, String).

Concorrência (anti-deadlock 40P01)

Dois arquivos do mesmo lote mesclam em paralelo (executor processamento, 4 threads), e ON CONFLICT DO UPDATE trava a linha conflitante mesmo quando o WHERE IS DISTINCT FROM não atualiza. Duas defesas, ambas obrigatórias:
  • ORDER BY pela chave natural no SELECT do INSERT — ordem de aquisição de row lock determinística entre transações concorrentes (convenção da frota bi-*; mesma correção do bi-comercial-xls).
  • Advisory lock de mescla (adquirirLockMescla(Consumer)) — serializa a fase de banco inteira; cobre as classes que o ORDER BY não cobre (UPSERT vs DELETE seletivo e DELETE vs DELETE), já que a ordem de lock do DELETE segue a ordem de scan e não é controlável.

Conexão

Obtida via ConnectionOperations.findConnectionStatus() para participar da transação ativa do chamador.
  • Constructor Details

    • BulkUpserter

      public BulkUpserter(io.micronaut.data.connection.ConnectionOperations<Connection> connectionOperations)
      Cria o upserter injetando o provedor de conexão JDBC usado em todo o fluxo idempotente.
      Parameters:
      connectionOperations - operações de conexão para participar da transação ativa do chamador (ver obterConexaoTransacional()).
  • Method Details

    • adquirirLockMescla

      public void adquirirLockMescla(Consumer<String> log)
      Serializa a fase de mescla entre processamentos concorrentes via pg_advisory_xact_lock: bloqueia até o lock ficar livre e o libera automaticamente no COMMIT/ROLLBACK da transação do chamador. Ver seção "Concorrência" no javadoc da classe.

      Exige transação ativa (falha alto via obterConexaoTransacional()); deliberadamente sem @Transactional — numa transação própria o lock seria liberado na saída do método, silenciosamente inútil. Chamar como primeira instrução de ProcessamentoTransactionalHelper.mesclarDadosDoPeriodo(LocalDate, LocalDate, List, List, Consumer).

      Parameters:
      log - callback de log textual — avisa quando houve espera relevante pelo lock.
    • mesclarFaturamentos

      public UpsertResultado mesclarFaturamentos(List<Faturamento> faturamentos, LocalDate periodoInicio, LocalDate periodoFim, Consumer<String> log)
      Mescla faturamentos em bi_faturamento preservando identidade de linhas inalteradas. Ver javadoc da classe para detalhes do fluxo.
      Parameters:
      faturamentos - lote de faturamentos do upload (pode ser vazio).
      periodoInicio - limite inferior (inclusivo) do dt_em para o DELETE seletivo.
      periodoFim - limite superior (inclusivo) do dt_em para o DELETE seletivo.
      log - callback de log textual (pode ser ignorado).
      Returns:
      UpsertResultado com contagens de inseridas/atualizadas/removidas e enviadas == faturamentos.size().
    • mesclarMovimentos

      public UpsertResultado mesclarMovimentos(List<Movimento> movimentos, LocalDate periodoInicio, LocalDate periodoFim, Consumer<String> log)
      Mescla movimentos em bi_movimento preservando identidade de linhas inalteradas. Antes do COPY, atribui seq_dentro_grupo via Map+AtomicInteger sobre as 16 cols do "evento físico" (spec § 4.1.b) — passa a ser a 17ª col da chave natural.
      Parameters:
      movimentos - lote de movimentos do upload (pode ser vazio).
      periodoInicio - limite inferior (inclusivo) do data_lcto para o DELETE seletivo.
      periodoFim - limite superior (inclusivo) do data_lcto para o DELETE seletivo.
      log - callback de log textual (pode ser ignorado).
      Returns:
      UpsertResultado com contagens de inseridas/atualizadas/removidas.
    • atribuirSeqDentroGrupo

      public static void atribuirSeqDentroGrupo(List<Movimento> movimentos)
      Atribui seq_dentro_grupo (1, 2, 3, ...) por ordem de iteração, indexado pela chave de 16 cols (ColunasConteudo.MOVIMENTO_CHAVE_GRUPO). Mutação in-place: sempre sobrescreve o valor atual para garantir consistência entre múltiplas chamadas no mesmo upload (spec § 4.1.b).

      Usado por mesclarMovimentos(List, LocalDate, LocalDate, Consumer) antes do COPY. Exposto como public static para facilitar testes que precisam reproduzir o cálculo (§ 9.7 da spec).

      Parameters:
      movimentos - lote de movimentos a numerar (mutado in-place).