Reserva dos emails da régua de trial em bloco

TLDR: a reserva dos emails da régua de trial passa a ser feita em blocos de 200 por query, com ON CONFLICT DO NOTHING RETURNING, em vez de um get_or_create por email. No pior caso de um lote (200 user trials × 6 emails), cai de ~1200 idas ao banco para 6.

Contexto

send_trial_journey_emails_batch processa um lote de até 200 user trials (JOURNEY_BATCH_SIZE). Para cada user trial, _dispatch_pending_emails percorre os emails devidos e chama _reserve_dispatch, que faz um get_or_create por email antes de enfileirar.

A reserva existe para garantir envio exatamente-uma-vez: o created do get_or_create é o que diz se esta execução ganhou a linha na constraint unique_together ("user", "source", "email_key"). Só quem ganha enfileira. Isso protege o job dos dois modos de falha reais do ambiente:

  • autoretry_for = (Exception,) em BaseTaskWithRetry reexecuta o lote inteiro do começo;
  • acks_late = True com visibility_timeout = 300 no SQS pode reentregar o mesmo lote a um segundo worker.

O custo é uma ida ao banco por email. Com 6 emails na régua, o teto de um lote é ~1200 reservas sequenciais. Esse teto só é atingido quando todos os user trials do lote estão sem nenhum dispatch registrado (primeira execução após deploy, backfill); em regime cada user trial cruza um marco por dia.

Medição do lote no pior caso: ~5s de banco contra ~30s de publicação no SQS (uma chamada HTTPS por email). O banco não é o gargalo — a mudança é de folga para crescimento, não de dor atual.

Objetivos

  • Reservar em bloco: uma query por bloco de 200 reservas, em vez de uma por email.
  • Preservar a garantia de exatamente-uma-vez: continuar sabendo, por reserva, se foi esta execução que a criou.
  • Preservar a ordem write-ahead: nenhuma ficha entra na fila sem a reserva já gravada.
  • Manter o retorno de send_trial_journey_emails_batch (status, user_trials, sent) inalterado.

Fora de escopo

  • Não altera JOURNEY_BATCH_SIZE (segue 200 user trials por task).
  • Não altera a régua, os templates, os percentuais nem resolve_pending_emails.
  • Não reduz as publicações no SQS — segue uma chamada por email.
  • Não adiciona soft_time_limit na task.
  • Não corrige a interação entre soft delete e a unique_together de EmailDispatch.

Decisão

bulk_create(ignore_conflicts=True) não serve. Com OnConflict.IGNORE o Django desliga o RETURNING (django/db/models/query.py:1940-1946), nenhuma PK é atribuída e o método devolve a mesma lista recebida — o retorno é idêntico tenha inserido tudo ou nada. Sem saber quais reservas foram desta execução, o retry reenfileira o lote inteiro e reenvia email já enviado.

update_conflicts=True devolve PKs, mas gera DO UPDATE SET: sobrescreveria o sent_at de reserva feita por outro worker e também não separa inserido de atualizado.

A combinação necessária — inserir em lote ignorando conflitos e saber quais linhas entraram — é INSERT ... ON CONFLICT DO NOTHING RETURNING, que o ORM não expressa. Por isso a reserva usa SQL parametrizado via connections["default"], exceção justificada à regra de preferir o ORM.

A escrita vai explicitamente em default porque DbRouter.db_for_read aponta para readonly (config/db_route.py:3); a reserva não pode ler réplica.

Bloco de 200

O bloco delimita a janela de órfãos: reservas gravadas cujo .delay() não chegou a acontecer porque a task morreu no meio da etapa de enfileiramento. Órfão nunca é reenviado, porque a reserva já existe e resolve_pending_emails o filtra via sent_keys — é perda silenciosa.

O teto técnico do Postgres é ~16.300 reservas por comando (65.535 valores / 4 por linha), então 200 é escolha de risco, não limite do banco. De 1200 idas para 6 já captura praticamente todo o ganho; blocos maiores reduziriam idas de forma irrelevante e aumentariam a perda máxima.

Falha durante a query de reserva não deixa órfão: a query falha inteira, nada é reservado e o retry refaz o bloco corretamente. A janela existe só entre a reserva confirmada e o fim do loop de enfileiramento daquele bloco.

Mudanças

  • apps/trials/tasks.py
    • remove _reserve_dispatch e _dispatch_pending_emails;
    • adiciona _pending_reservations, que monta todos os pares (user_id, email_key) devidos do lote sem tocar o banco (sent_keys_by_user já carrega o histórico em uma query);
    • adiciona _reserve_block, que insere um bloco com ON CONFLICT DO NOTHING RETURNING user_id, email_key e devolve o conjunto efetivamente inserido;
    • send_trial_journey_emails_batch passa a iterar os pendentes em blocos de RESERVATION_BLOCK_SIZE, reservando o bloco e enfileirando apenas os pares retornados.
  • tests/trials/test_tasks.py — cobre reserva em bloco, contagem de queries, não-reenvio no retry, e que só o vencedor da constraint enfileira.

Como testa

bash make run.test path=tests/trials/test_tasks.py

Cenários cobertos:

  • lote com vários emails devidos reserva em uma query e enfileira todos;
  • segunda execução do mesmo lote não enfileira nada (reservas já existem);
  • pendentes acima do tamanho do bloco geram uma query por bloco;
  • reserva concorrente já gravada por outra execução é pulada, sem reenvio;
  • SEND_EMAIL=False não reserva nem enfileira.