Depois de ler este artigo você vai saber mover trabalho lento para fora da requisição HTTP usando filas e workers, e vai conseguir deixar esse trabalho seguro contra repetições, falhas e picos. Os exemplos usam Celery (Python) e BullMQ (Node), mas os conceitos valem para RabbitMQ, SQS, Kafka e afins.
O problema real
O usuário envia um CSV grande. A requisição valida o arquivo, grava no banco, manda um e-mail e gera um PDF. Tudo no mesmo ciclo HTTP, e o navegador desiste antes de terminar. O usuário acha que deu erro e envia de novo, e agora você processa o arquivo duas vezes.
A correção é separar o que precisa acontecer agora do que pode acontecer depois. A requisição valida rápido, grava o arquivo, entrega um trabalho à fila e responde com 202 Accepted e um identificador. Um worker, que é um processo separado, pega o trabalho e executa.
# Depois: a requisição só enfileira
@app.post("/upload-csv")
def upload_csv(arquivo):
caminho = salvar_em_storage(arquivo) # S3, disco compartilhado etc.
tarefa = processar_csv.delay(caminho) # passa a referência, não o conteúdo
return {"status": "processando", "tarefa_id": tarefa.id}, 202
Passe sempre identificadores ou caminhos para a tarefa, nunca o conteúdo inteiro. Mensagens grandes sobrecarregam o broker e ficam desatualizadas.
Produtor, broker e consumidor
O produtor coloca mensagens. O broker guarda até alguém consumir. O consumidor (worker) processa e confirma.
Celery com Redis
# tasks.py
from celery import Celery
app = Celery("tasks", broker="redis://localhost:6379/0")
@app.task(bind=True, acks_late=True, time_limit=300)
def processar_csv(self, caminho: str):
... # lógica pesada
return {"status": "ok"}
Para subir o worker: celery -A tasks worker --loglevel=info. A opção acks_late=True só confirma a mensagem depois que a tarefa termina, então um worker que morre no meio devolve o trabalho à fila. O preço é que a tarefa pode rodar mais de uma vez, o que nos leva à idempotência.
BullMQ com Redis
O pacote bull está em modo de manutenção. Para projetos novos em Node use o bullmq:
import { Queue, Worker } from "bullmq";
const connection = { host: "localhost", port: 6379 };
const fila = new Queue("csv", { connection });
// Produtor
app.post("/upload", async (req, res) => {
const job = await fila.add(
"processar",
{ caminho: req.body.caminho },
{ attempts: 5, backoff: { type: "exponential", delay: 2000 } }
);
res.status(202).json({ jobId: job.id });
});
// Consumidor (processo separado)
new Worker("csv", async (job) => {
console.log("processando", job.data.caminho);
}, { connection, concurrency: 5 });
Garantias de entrega
| Garantia | Significado | Quando serve |
|---|---|---|
| At-most-once | Roda no máximo uma vez, pode não rodar | Notificação descartável |
| At-least-once | Roda pelo menos uma vez, pode repetir | Quase tudo |
| Exactly-once | Roda exatamente uma vez | Só dentro de sistemas específicos |
Entre dois sistemas separados, exactly-once de ponta a ponta não existe de graça. O combinado que funciona é entrega at-least-once somada a tarefas idempotentes.
Idempotência
Uma tarefa é idempotente quando rodar duas vezes deixa o mesmo resultado que rodar uma. A técnica mais comum é uma chave de idempotência com restrição de unicidade no banco:
def registrar_transferencia(transfer_id: str, origem: int, destino: int, valor: int):
with db.transaction():
inserido = db.execute(
"INSERT INTO transferencias (id, origem, destino, valor) "
"VALUES (%s, %s, %s, %s) ON CONFLICT (id) DO NOTHING",
(transfer_id, origem, destino, valor),
).rowcount
if inserido == 0:
return # já processada, ignora
db.execute("UPDATE contas SET saldo = saldo - %s WHERE id = %s", (valor, origem))
db.execute("UPDATE contas SET saldo = saldo + %s WHERE id = %s", (valor, destino))
O transfer_id vem de quem originou a operação. Se a mensagem chegar duas vezes, o segundo INSERT não insere nada e a tarefa sai sem mexer nos saldos. O registro e os dois UPDATE ficam na mesma transação, então ou tudo acontece ou nada.
Retry com backoff exponencial e jitter
Repetir imediatamente uma chamada que falhou costuma piorar o problema, porque todos os trabalhos falhos voltam juntos. O certo é esperar um tempo crescente, com um componente aleatório (jitter) para espalhar as tentativas.
No Celery há suporte nativo:
@app.task(
bind=True,
autoretry_for=(ConnectionError, TimeoutError),
retry_backoff=True, # 1s, 2s, 4s, 8s...
retry_backoff_max=300, # teto de 5 minutos
retry_jitter=True, # aleatoriza cada espera
max_retries=5,
)
def enviar_email(self, email_id: int):
provedor.enviar(email_id)
Repare em autoretry_for: só vale a pena repetir erros transitórios. Um erro de validação vai falhar sempre, e repetir só desperdiça recurso.
Dead letter queue
Depois de esgotar as tentativas, o trabalho precisa ir para algum lugar onde alguém possa olhar, em vez de sumir ou ficar repetindo. Essa é a dead letter queue (DLQ).
O Celery não tem uma DLQ pronta. Com RabbitMQ você configura um dead letter exchange na fila, e as mensagens rejeitadas ou expiradas vão para outra fila. Com SQS, configure uma redrive policy com maxReceiveCount. Com o BullMQ, os jobs que falharam ficam no conjunto de falhos e podem ser inspecionados e reenfileirados. Em todos os casos, crie um alerta para quando a DLQ deixar de estar vazia.
Ordem versus paralelismo
Mais workers processam mais rápido, mas sem garantia de ordem entre mensagens. Nem toda tarefa liga para isso: enviar mil e-mails em qualquer ordem é inofensivo.
Quando a ordem importa só dentro de um grupo, por exemplo todas as alterações de um mesmo cliente, particione pela chave do grupo. Em Kafka, isso é a chave da mensagem, que determina a partição, e cada partição é consumida por um consumidor do grupo por vez. No RabbitMQ, uma opção é uma fila por grupo ou o plugin de consistent hash exchange. Se você precisa de ordem global estrita, aceite que vai processar um de cada vez.
Monitoramento: profundidade e atraso
Dois números contam quase toda a história. O tamanho da fila e há quanto tempo a tarefa mais antiga espera.
Com o broker em Redis, o Celery guarda cada fila numa lista:
redis-cli LLEN celery
Para o atraso, o jeito confiável é a própria tarefa registrar o horário de enfileiramento e o worker calcular agora - enfileirado_em ao começar a executar, publicando isso como histograma. Alertas úteis ficam em torno de "o atraso passou do aceitável para o negócio por N minutos", e não em um tamanho fixo de fila, pois 10 mil tarefas curtas podem ser nada e 50 tarefas longas podem ser um problema. Veja como expor e alertar métricas no artigo sobre monitoramento com Prometheus e Grafana.
Aplicação em IA: inferência longa fora da requisição
Chamadas a um LLM com saída longa ou a um modelo local podem levar dezenas de segundos. O padrão é o mesmo: enfileirar, devolver um job_id e deixar o cliente consultar o resultado, ou receber por webhook ou SSE.
@app.post("/analisar")
def analisar(texto: str):
job = analisar_texto.delay(texto)
return {"job_id": job.id}, 202
@app.get("/resultado/{job_id}")
def resultado(job_id: str):
r = AsyncResult(job_id, app=celery_app)
if r.ready():
return {"status": "pronto", "resultado": r.get()}
return {"status": "processando"}
Para isso funcionar, configure um backend de resultados no Celery (backend="redis://...") e um prazo de expiração para os resultados.
Armadilhas comuns
- Fila sem limite de crescimento e Redis sem
maxmemoryadequado. Se o mesmo Redis serve de cache e de broker, uma política de remoção comoallkeys-lrupode apagar mensagens. Use instâncias separadas. - Tarefa sem
time_limit, que fica presa para sempre esperando uma API externa. - Retry infinito para erro que nunca se resolve.
- Efeitos colaterais não idempotentes, como e-mail enviado duas vezes sem nenhuma chave de controle.
Quando não usar fila
Se a tarefa leva poucos milissegundos e o usuário precisa do resultado agora, a fila só acrescenta latência e complexidade. Se você ainda não tem um broker na infraestrutura, comece com algo simples como uma tabela de jobs no Postgres consumida com SELECT ... FOR UPDATE SKIP LOCKED, e migre quando a carga pedir.
Próximos passos
Para usar o Redis como cache sem misturar com o broker, leia Cache com Redis na prática. Para provisionar broker e workers como código, veja Terraform: primeiros passos.
Para testar, suba um Redis com Docker, crie uma tarefa que dorme 5 segundos, enfileire cem de uma vez e observe quantos workers e qual concurrency dão o melhor tempo total.
Quer aplicar isso na sua empresa? Marque uma conversa de 45 minutos em https://iauaicloud.com.br/consultoria