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-progressde 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-progresspara 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:
- Grava a última posição em
progress:latest_position:{content_type}:{content_id}:{user_id}(TTL 1h) - Adiciona
{content_id}:{user_id}ao setprogress:pending_flush:{content_type}(SADD) - 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:
- Para cada
content_type, drena o setprogress:pending_flush:{content_type}viapop_pending_batches(SPOPem lotes) - Fan-out de tasks
save_progress_batch(content_type, content_user_ids)com atéPROGRESS_FLUSH_BATCH_SIZEentradas cada - Cada
save_progress_batchlê a última posição de cadacontent_user_ide chama o{Lesson,Audio,Meditation,Live}ProgressService.save_progressexistente — toda a lógica de negócio é reusada como está - Em falha por entrada,
requeue_pendingdevolve ocontent_user_idao 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
/progressoverview/summary leem dela; hoje o atraso depende da profundidade da fila e é pior sob carga) - Completion e
advance_user_journeydisparam 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 deREDIS_URLapps/engagements/services/progress_buffer.py(novo) —ProgressBufferServiceapps/engagements/tasks.py— adicionarflush_pending_progressesave_progress_batch(filaonion-progress); manter assave_progress_*legadasapps/engagements/serializers.py—ProgressLessonSerializer.applyusaProgressBufferServiceem 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 liveconfig/celery_defaults.py— entrada no beat schedule (PROGRESS_FLUSH_INTERVAL_SECONDS, default 60)config/settings/base.py— env varsPROGRESS_FLUSH_INTERVAL_SECONDSePROGRESS_FLUSH_BATCH_SIZE(default 500)- Testes em
tests/engagements/— unit tests doProgressBufferService(mock de redis), testes deflush_pending_progress/save_progress_batche dos serializers atualizados
Como verificar
- Unit tests: store/pop do buffer, fan-out de lotes, requeue por entrada em falha, serializers não enfileiram mais
- Local: fazer POST em
/audios/progresse/engagements/progress/lessonrepetidamente para o mesmo usuário/conteúdo → confirmar um únicoUPDATEno banco por janela de flush, com a última posição - Completion: POST com posição near-end → após o próximo flush,
completed_até setado eadvance_user_journeyé enfileirada exatamente uma vez - Sobreposição de deploy: publicar manualmente uma task legada
save_progress_audio→ ainda processa - Staging: comparar a taxa de mensagens da SQS
onion-progressantes/depois sob carga simulada
Documentação
- learnings/infrastructure_sqs_countdown_visibility_timeout.md — Celery ETA/countdown ≥
visibility_timeoutdo SQS causa execução duplicada; usar flush periódico via beat no lugar