Progresso com buffer: Redis + flush periódico

TLDR: Parar de enfileirar uma task Celery por heartbeat de progresso (a cada 5s por usuário ativo); guardar a última posição no Redis e descarregar no banco com uma task periódica, reduzindo o tráfego da fila onion-progress de O(usuários ativos × 12/min) para poucas tasks por minuto.

Contexto

O app envia heartbeats de progresso a cada ~5s enquanto o usuário assiste uma aula ou ouve um áudio:

  • Aulas: o player atualiza a cada 1s, com gate client-side de 1 save por 5s → POST /engagements/progress/lesson
  • Áudio: polling de 5s com checagem de movimento de posição ≥1s → POST /audios/progress
  • Meditação e live seguem o mesmo padrão

No backend, cada POST enfileira imediatamente uma task Celery (save_progress_*, apps/engagements/tasks.py) na fila onion-progress, e cada task faz um UPDATE no banco. Com N ouvintes simultâneos isso produz N × 12 tasks/min, inundando a fila mesmo depois de ela ter sido isolada (ver 20260604120000_dedicated_queues_checkout_progress.md).

A correção precisa ficar no backend: mudança no app exige ciclo de release na loja. A análise do app confirmou que as duas chamadas são fire-and-forget — o corpo da resposta ({"task_id": ...}) nunca é lido e erros são engolidos silenciosamente — então o request path pode mudar livremente desde que retorne 201.

Só a última posição importa por (usuário, conteúdo): heartbeats intermediários são redundantes. A completion é derivada no servidor (position >= duration - 10s), então fica preservada desde que a posição final eventualmente chegue ao banco.

Por que não debounce por usuário com countdown

O broker de produção é SQS com visibility_timeout: 30 (config/celery_settings.py). Tasks com ETA/countdown no SQS ficam retidas sem ack no worker; um countdown ≥ visibility timeout causa reentrega da mensagem e execução duplicada. Além disso, continuaria produzindo O(usuários ativos) tasks. Rejeitado.

Por que não throttling DRF (429)

Descartar requests arrisca perder a posição final de uma sessão (inclusive a que dispara a completion). Rejeitado.

Objetivos

  • Reduzir o tráfego da fila onion-progress para uma taxa constante, independente do número de usuários
  • Reduzir escritas no banco de 1 por heartbeat para no máximo 1 por (usuário, conteúdo) por janela de flush
  • Nunca perder a posição final de uma sessão; detecção de completion inalterada
  • Sem mudanças no app; contrato da API continua POST → 201
  • Sem mudanças nas regras de negócio de progresso (completion near-end, advance_user_journey)

Fora de escopo

Aumentar o intervalo de heartbeat para 15–30s e enviar flush em eventos de pause/background — reduz ainda mais o tráfego HTTP, mas não é necessário para a correção da fila (follow-up futuro no app).

Mudanças

Nomenclatura

Item Nome Significado
Chave Redis (valor) progress:latest_position:{content_type}:{content_id}:{user_id} última posição reportada daquele conteúdo para aquele usuário
Chave Redis (set) progress:pending_flush:{content_type} conjunto de content_user_ids ("{content_id}:{user_id}") aguardando persistência
Service ProgressBufferService store_latest_position, pop_pending_batches, get_latest_position, requeue_pending
Task de beat flush_pending_progress varre os sets de pendências e despacha os lotes
Task de lote save_progress_batch(content_type, content_user_ids) persiste um lote, na mesma família das save_progress_*
Env vars PROGRESS_FLUSH_INTERVAL_SECONDS, PROGRESS_FLUSH_BATCH_SIZE janela do flush (60) e tamanho do lote (500)

content_type ∈ lesson, audio, meditation, live.

Request path (síncrono, sem Celery)

O apply() dos serializers de progresso para de chamar .delay(). Em vez disso:

  1. Grava a última posição em progress:latest_position:{content_type}:{content_id}:{user_id} (TTL 1h)
  2. Adiciona {content_id}:{user_id} ao set progress:pending_flush:{content_type} (SADD)
  3. Retorna 201 com {"data": true}

Duas operações Redis O(1) substituem um publish no SQS por request — o request também fica mais rápido.

O RedisCache nativo do Django não expõe operações de SET, então um pequeno client raw compartilhado (redis-py a partir de settings.REDIS_URL) é introduzido em apps/common/.

Flush path (periódico, taxa constante de tasks)

Uma nova entrada de beat roda flush_pending_progress a cada 60s na onion-progress:

  1. Para cada content_type, drena o set progress:pending_flush:{content_type} via pop_pending_batches (SPOP em lotes)
  2. Fan-out de tasks save_progress_batch(content_type, content_user_ids) com até PROGRESS_FLUSH_BATCH_SIZE entradas cada
  3. Cada save_progress_batch lê a última posição de cada content_user_id e chama o {Lesson,Audio,Meditation,Live}ProgressService.save_progress existente — toda a lógica de negócio é reusada como está
  4. Em falha por entrada, requeue_pending devolve o content_user_id ao set para o próximo flush reprocessar

Lotes pequenos mantêm cada task bem abaixo dos 30s de visibility timeout do SQS; todas as escritas são idempotentes, então uma reentrega rara é inofensiva.

Carga resultante

  Antes Depois
Tasks/min (5k ouvintes) ~60.000 ~1 flush + ~10 lotes
Escritas no banco/min (5k ouvintes) ~60.000 ≤ 5.000 (1 por usuário a cada 60s)

Trade-offs (aceitos)

  • A posição no banco atrasa até ~65s (as telas de /progress overview/summary leem dela; hoje o atraso depende da profundidade da fila e é pior sob carga)
  • Completion e advance_user_journey disparam até ~65s depois de o usuário atingir o near-end
  • Se o Redis perder uma entrada entre heartbeats e flush, no máximo ~60s de posição são perdidos — o próximo heartbeat restaura

Segurança no rollout

As definições das tasks save_progress_* são mantidas (mensagens SQS em trânsito durante o deploy continuam processando); só os pontos de enfileiramento mudam. Remoção em release posterior.

Arquivos

  • apps/common/redis.py (novo) — client raw redis-py compartilhado a partir de REDIS_URL
  • apps/engagements/services/progress_buffer.py (novo) — ProgressBufferService
  • apps/engagements/tasks.py — adicionar flush_pending_progress e save_progress_batch (fila onion-progress); manter as save_progress_* legadas
  • apps/engagements/serializers.py — ProgressLessonSerializer.apply usa ProgressBufferService em vez de .delay()
  • apps/audios/serializers/progress_audio.py, apps/audios/serializers/progress_meditation.py, apps/lives/serializers.py — idem para áudio, meditação e live
  • config/celery_defaults.py — entrada no beat schedule (PROGRESS_FLUSH_INTERVAL_SECONDS, default 60)
  • config/settings/base.py — env vars PROGRESS_FLUSH_INTERVAL_SECONDS e PROGRESS_FLUSH_BATCH_SIZE (default 500)
  • Testes em tests/engagements/ — unit tests do ProgressBufferService (mock de redis), testes de flush_pending_progress/save_progress_batch e dos serializers atualizados

Como verificar

  1. Unit tests: store/pop do buffer, fan-out de lotes, requeue por entrada em falha, serializers não enfileiram mais
  2. Local: fazer POST em /audios/progress e /engagements/progress/lesson repetidamente para o mesmo usuário/conteúdo → confirmar um único UPDATE no banco por janela de flush, com a última posição
  3. Completion: POST com posição near-end → após o próximo flush, completed_at é setado e advance_user_journey é enfileirada exatamente uma vez
  4. Sobreposição de deploy: publicar manualmente uma task legada save_progress_audio → ainda processa
  5. Staging: comparar a taxa de mensagens da SQS onion-progress antes/depois sob carga simulada

Documentação