Queues/Workers/async agents
Você construirá um worker recuperável para criar tickets sintéticos. Comece pelo banco local persistente e depois use Redis da semana 42 para examinar entrega e confirmação. A entrega deve mostrar que reinício não elimina job e que repetição não duplica o efeito. Falha permanente precisa aparecer em dead com um caminho de revisão.
O script principal requer Python 3 e usa somente sqlite3. Use jobs-lab.sqlite em uma pasta exclusiva; não conecte a bancos de produção. A etapa de Redis requer a stack Docker anterior. Caso você escolha Celery como extensão, registre versão, broker e configuração, e não trate o exemplo SQLite como evidência de execução de Celery.
PythonAgentesInfraestruturaAo terminar esta aula
- Fila, estado de job e efeito são contratos diferentes.
- Idempotência precisa sobreviver a reinício.
- Retries têm causa, espera, orçamento e limite.
Antes de continuar: Leitura: Queues/Workers/async agents
Definir estados e contrato de consulta
FundamentosEscreva uma tabela de transições permitidas: pending para running, running para succeeded, retry_wait ou dead, e retry_wait para pending. Defina erro definitivo e temporário. Acrescente tenant e operationHash ao job e exija a identidade na consulta de status. Um pedido de criação só retorna aceitação após commit. Teste inserção duplicada e chave com argumentos diferentes. A primeira pode recuperar o job existente; a segunda deve ser conflito. Separe status técnico de conclusão funcional: um job pode terminar com encaminhamento correto sem ter criado um estorno.
Executar persistência e repetir com segurança
SegurançaRode jobs.py, registre as duas linhas e a contagem de ticket. Feche o processo e execute novamente com o mesmo arquivo. j1 permanece succeeded e continua com um ticket. Faça um teste com conexão nova para comprovar que o resultado vem do banco, não de variável em memória. Acrescente uma exceção entre a inserção do ticket e o update do job dentro da transação; o rollback deve impedir estado parcial. Depois mova a inserção para fora da transação em uma cópia descartável e observe a janela de inconsistência que você introduziu.
import sqlite3
db=sqlite3.connect('jobs-lab.sqlite')
db.executescript('''
CREATE TABLE IF NOT EXISTS jobs(id TEXT PRIMARY KEY,state TEXT,attempt INTEGER);
CREATE TABLE IF NOT EXISTS tickets(job_id TEXT PRIMARY KEY,subject TEXT);
''')
with db:
db.execute("INSERT OR IGNORE INTO jobs VALUES ('j1','pending',0)")
db.execute("INSERT OR IGNORE INTO jobs VALUES ('j2','pending',0)")
def work(job_id, subject):
with db:
row=db.execute('SELECT state FROM jobs WHERE id=?',(job_id,)).fetchone()
if not row: raise ValueError('unknown-job')
if row[0]=='succeeded': return 'already-done'
if row[0]=='dead': return 'needs-review'
db.execute('UPDATE jobs SET attempt=attempt+1 WHERE id=?',(job_id,))
if not subject:
db.execute("UPDATE jobs SET state='dead' WHERE id=?",(job_id,))
return 'invalid-input'
db.execute('INSERT OR IGNORE INTO tickets VALUES (?,?)',(job_id,subject))
db.execute("UPDATE jobs SET state='succeeded' WHERE id=?",(job_id,))
return 'done'
work('j1','consulta');work('j1','consulta');work('j2','')
assert db.execute("SELECT count(*) FROM tickets WHERE job_id='j1'").fetchone()[0]==1
assert db.execute("SELECT state FROM jobs WHERE id='j2'").fetchone()[0]=='dead'
print(db.execute('SELECT * FROM jobs ORDER BY id').fetchall())
db.close()Modelar retry com relógio e orçamento
FundamentosAdicione next_attempt_at, last_error e max_attempts. Use um relógio injetável e uma dependência falsa que falha duas vezes com erro temporário antes de funcionar. O worker só deve retomar quando a espera vencer, e o terceiro sucesso não pode gerar mais de um ticket. Para argumento inválido, entre em dead imediatamente. Teste deadline vencido e falta de orçamento antes da próxima tentativa. Conte todas as tentativas no relatório, distinguindo chamadas ao modelo e ao efeito. Não use loops de retry sem espera e sem limite.
Verificar Redis Streams com confirmação
RedisNa stack anterior, use redis-cli XADD lab:jobs * jobId j3, depois XGROUP CREATE lab:jobs lab-workers 0 MKSTREAM. Consuma com XREADGROUP GROUP lab-workers worker-1 COUNT 1 STREAMS lab:jobs >; em shell, coloque > entre aspas para não redirecionar saída. Consulte XPENDING para ver a entrada pendente e confirme com XACK usando o ID recebido. Antes do ack, simule perda do consumidor e examine o mecanismo de recuperação suportado na versão instalada. Esses comandos exercitam o broker real, mas o efeito precisa continuar protegido no banco.
Implementar dead e demonstrar redrive consciente
FundamentosInsira um job com entrada inválida, confirme dead e registre motivo sanitizado. Corrija a entrada somente após revisão e preserve a identidade da operação, com hash atualizado por um procedimento explícito. Faça redrive e confirme um efeito, não dois. Entregue evidência de reinício, retry temporário, falha permanente e recuperação de mensagem pendente. Descreva a fronteira entre gravar job e publicar em Redis, propondo outbox para evitar perda entre as duas operações. A conclusão deve indicar quais garantias foram testadas e quais dependem de lease, concorrência e idempotência do serviço externo.
Exercício aplicado
O worker cria um ticket e cai antes de confirmar a mensagem. Outro worker recebe o mesmo job. Explique a recuperação segura e por que ack tardio não basta.
- Busque o estado e o resultado por chave durável.
- Verifique unicidade e assinatura da operação.
- Recupere o ticket existente em vez de recriá-lo.
- Confirme a mensagem após concluir o protocolo.
Abrir resolução comentada
Redelivery é esperado em desenhos que priorizam recuperação. A chave de operação e o resultado durável permitem reconhecer o efeito anterior. A confirmação da fila informa transporte; ela não desfaz nem deduplica o ticket automaticamente.
No exemplo, efeito e estado estão no mesmo banco e podem participar de uma transação. Para API externa, é necessário um protocolo de idempotência e reconciliação com o destino. Retry e dead-letter também precisam preservar a identidade lógica do trabalho.
import sqlite3
db=sqlite3.connect('jobs-lab.sqlite')
db.executescript('''
CREATE TABLE IF NOT EXISTS jobs(id TEXT PRIMARY KEY,state TEXT,attempt INTEGER);
CREATE TABLE IF NOT EXISTS tickets(job_id TEXT PRIMARY KEY,subject TEXT);
''')
with db:
db.execute("INSERT OR IGNORE INTO jobs VALUES ('j1','pending',0)")
db.execute("INSERT OR IGNORE INTO jobs VALUES ('j2','pending',0)")
def work(job_id, subject):
with db:
row=db.execute('SELECT state FROM jobs WHERE id=?',(job_id,)).fetchone()
if not row: raise ValueError('unknown-job')
if row[0]=='succeeded': return 'already-done'
if row[0]=='dead': return 'needs-review'
db.execute('UPDATE jobs SET attempt=attempt+1 WHERE id=?',(job_id,))
if not subject:
db.execute("UPDATE jobs SET state='dead' WHERE id=?",(job_id,))
return 'invalid-input'
db.execute('INSERT OR IGNORE INTO tickets VALUES (?,?)',(job_id,subject))
db.execute("UPDATE jobs SET state='succeeded' WHERE id=?",(job_id,))
return 'done'
work('j1','consulta');work('j1','consulta');work('j2','')
assert db.execute("SELECT count(*) FROM tickets WHERE job_id='j1'").fetchone()[0]==1
assert db.execute("SELECT state FROM jobs WHERE id='j2'").fetchone()[0]=='dead'
print(db.execute('SELECT * FROM jobs ORDER BY id').fetchall())
db.close()Como conferir seu resultado
- Reexecução após reinício mantém um ticket.
- Falha dentro da transação não deixa efeito parcial.
- Erro permanente não entra em retry infinito.
- Entrada pendente é inspecionada e confirmada no Redis.
Aplique em um problema novo
Primeiro resolva sem consultar a resposta. Explique suas decisões e guarde a evidência. A conclusão de leitura é independente desta autoavaliação.
Confira seus pré-requisitos
- Distinguir ack e efeito do job.
- Usar identidade/ledger duráveis.
Job K2 cria nota e grava resultado em transação; worker cai antes de ack. Segunda entrega K2 chega. Job K3 tem argumento inválido e vai para DLQ. Supervisor reenvia K3 após corrigir dados.
Conferir raciocínio e critérios de domínio
K2 replay consulta ledger e devolve recibo; total um efeito. Ack da segunda entrega confirma consumo sem nova nota.
K3 inválido não recebe retries cegos; DLQconserva motivo/contexto sanitizado para revisão.
Reenvio supervisionado precisa preservar ou versionar intenção conforme mudança dos dados. Se payload mudou sob mesma chave, rejeitar incompatibilidade em vez de aceitar silenciosamente.
Evidências para autoavaliação ou revisão por pares
- Replay e ack: K2 replay consulta ledger e devolve recibo; total um efeito. Ack da segunda entrega confirma consumo sem nova nota.
- Erro permanente: K3 inválido não recebe retries cegos; DLQconserva motivo/contexto sanitizado para revisão.
- Payload no reenvio: Reenvio supervisionado precisa preservar ou versionar intenção conforme mudança dos dados. Se payload mudou sob mesma chave, rejeitar incompatibilidade em vez de aceitar silenciosamente.
Um erro frequente
DLQ elimina necessidade de investigar.
DLQ conserva falhas para recuperação supervisionada.
Teste sua compreensão
Responda com suas palavras antes de abrir o comentário. Saber explicar uma decisão é parte do domínio.
1. Redelivery implica necessariamente efeito duplicado?
2. Dead-letter conclui a tarefa do usuário?
Não.
Ela separa trabalho problemático para revisão e recuperação.
3. Por que um outbox é útil?
Para registrar estado e intenção de publicação juntos.
Publicação independente do commit cria janelas de perda ou inconsistência.
Seu progresso fica salvo neste navegador. Concluir a leitura não substitui demonstrar o domínio nos exercícios.
Referências e aprofundamento
Documentação oficial e trabalhos originais. As referências registram o escopo e as limitações para você conferir o que sustentam.
- Redis streaming
Redis • consulta: 2026-10-06
RedisAPIsStreams, logs append-only, consumer groups, acknowledgments e distribuição entre workers.
Limites: Entrega at-least-once requer tratamento de duplicações; Redis Streams e uma fila de jobs têm semânticas diferentes.
- Tasks
Celery • consulta: 2026-10-06
FundamentosTasks, workers, execução assíncrona, idempotência, acknowledgment, retries e limites de tempo.
Limites: Retry e redelivery podem repetir efeitos; desenhar tarefa idempotente e testar interrupções.
- Asynchronous Request-Reply Pattern
Microsoft • consulta: 2026-10-06
FundamentosSeparar aceite de requisição, processamento longo e consulta de status.
Limites: Persistir status e tratar timeout/cancelamento; assíncrono não elimina falhas.
- Retry pattern
Microsoft • consulta: 2026-10-06
FundamentosRetries de falhas transitórias, política de atrasos, logging e impacto em transações.
Limites: Retries de operações não idempotentes e camadas aninhadas podem duplicar efeitos e ampliar carga.
- Using dead-letter queues in Amazon SQS
AWS • consulta: 2026-10-06
InfraestruturaDLQ, política de redrive, máximo de recebimentos e retenção.
Limites: DLQ exige inspeção e política de reprocessamento; seu uso pode interferir em ordenação estrita FIFO.