O problema real
Usuário faz upload de um CSV com 100 mil linhas. Você valida, insere no banco, envia email e gera relatório PDF. Tudo na requisição HTTP. Tempo total: 45 segundos. Browser timeout, usuário pensa que foi erro.
Realidade: validação e injeção de dados são rápidos. Email demora 5 segundos (rede pra provider), PDF demora 15 segundos (processamento). Colocar tudo na requisição é punição.
Solução: coloca validação na requisição, mas joga validação + envio + PDF pra fila. Responde imediatamente ("recebemos seu upload, processando..."), processa em background.
Quando tirar da requisição HTTP e jogar para fila
# ❌ ANTES: tudo na requisição (45 segundos, timeout)
@app.post("/upload-csv")
def upload_csv(file):
# Validação rápida
df = pd.read_csv(file)
# Inserção em banco (5 segundos)
for row in df.iterrows():
db.insert(row)
# Email (5 segundos)
send_email(user, "CSV importado")
# PDF (15 segundos)
pdf = generate_report(df)
return {"status": "done"} # responde após 45 segundos
# ✓ DEPOIS: validação na requisição, resto em fila
@app.post("/upload-csv")
def upload_csv(file):
# Validação rápida (< 1 segundo)
df = pd.read_csv(file)
# Joga pra fila
queue.enqueue(process_csv, file_id, df)
return {"status": "processing", "job_id": "xyz"} # responde em 1 segundo
def process_csv(file_id, df): # roda em worker, pode demorar
# Inserção em banco
for row in df.iterrows():
db.insert(row)
# Email
send_email(user, "CSV importado")
# PDF
pdf = generate_report(df)
Produtor/Consumidor com código funcional
Usaremos RabbitMQ (na prática) ou Celery (Python) ou BullMQ (Node):
Python com Celery + Redis
from celery import Celery
# Configuração
app = Celery('tasks', broker='redis://localhost:6379')
# PRODUTOR: define tarefa
@app.task
def process_csv(file_id, csv_data):
# Processa em background
print(f"Processando arquivo {file_id}")
# ... lógica pesada
return {"status": "done"}
# APLICAÇÃO FLASK
from flask import Flask, jsonify
flask_app = Flask(__name__)
@flask_app.route('/upload', methods=['POST'])
def upload():
# Coloca na fila
task = process_csv.delay(file_id=123, csv_data="...")
# Retorna imediato
return jsonify({"task_id": task.id, "status": "queued"}), 202
# CONSUMIDOR: worker que processa
# terminal: celery -A tasks worker --loglevel=info
Node com Bull + Redis
const Queue = require('bull');
const redis = require('redis');
// Criar fila
const csvQueue = new Queue('csv-processing', {
redis: { host: 'localhost', port: 6379 }
});
// PRODUTOR: adicionar job
app.post('/upload', async (req, res) => {
const job = await csvQueue.add(
{ fileId: 123, csvData: "..." },
{ attempts: 3, backoff: { type: 'exponential', delay: 2000 } }
);
res.status(202).json({ jobId: job.id, status: 'queued' });
});
// CONSUMIDOR: processar job
csvQueue.process(async (job) => {
console.log(`Processando arquivo ${job.data.fileId}`);
// ... lógica pesada
return { status: 'done' };
});
// Monitor progresso
csvQueue.on('completed', (job) => {
console.log(`Job ${job.id} completou`);
});
Garantias de entrega: at-most-once vs at-least-once
| Garantia | O que significa | Risco | Uso |
|---|---|---|---|
| At-most-once | Tarefa roda 0 ou 1 vez | Se falha, nunca roda | Notificação não-crítica |
| At-least-once | Tarefa roda 1+ vezes | Se falha, pode rodar 2x | Email, relatório |
| Exactly-once | Tarefa roda exatamente 1 | Complexo, overhead | Transação financeira |
Realidade: exactly-once é quase sempre marketing. Use at-least-once + idempotência.
IDEMPOTÊNCIA: o requisito não-negociável
Se tarefa roda 2 vezes, resultado é idêntico:
# ❌ NÃO-IDEMPOTENTE: roda 2x, saldo muda 2x
def transfer_money(from_account, to_account, amount):
from_account.balance -= amount
to_account.balance += amount
db.save() # Se falha e retry, transfer 2x!
# ✓ IDEMPOTENTE: roda N vezes, resultado igual
def transfer_money(transfer_id, from_account, to_account, amount):
# Chave de idempotência previne duplicação
existing = db.query(f"SELECT * FROM transfers WHERE id = {transfer_id}")
if existing:
return existing # já transferiu, retorna resultado anterior
from_account.balance -= amount
to_account.balance += amount
db.save(transfer_id)
# Cliente sempre manda transfer_id:
transfer_money(
transfer_id="uuid-abcd-1234",
from_account=123,
to_account=456,
amount=100
)
Sem idempotência, fila retry duplo = seu usuário transfere $100 duas vezes.
Retry com backoff exponencial + jitter
Não retry imediato (sobrecarrega sistema). Espera crescente:
# ❌ Retry imediato
@app.task(bind=True)
def send_email(self, email):
try:
provider.send(email)
except Exception as e:
self.retry(countdown=0) # retry imediato
# ✓ Retry com backoff exponencial + jitter
import random
@app.task(bind=True)
def send_email(self, email):
try:
provider.send(email)
except Exception as e:
# Tentativa 1: espera 1 segundo
# Tentativa 2: espera 2 segundos
# Tentativa 3: espera 4 segundos
# Tentativa 4: espera 8 segundos
# + jitter (±20%): 1.2s, 2.1s, 3.8s, 9.1s
countdown = (2 ** self.request.retries) + random.randint(0, 20)
self.retry(countdown=countdown, max_retries=4)
Resultado real:
- Retry imediato: 100 falhas = 100 requests em 1 segundo (DDOS seu próprio servidor)
- Retry com backoff: 100 falhas = distribuídas ao longo de 15 segundos (sistema respira)
Dead Letter Queue (DLQ)
Jobs que falham 3 vezes vão pra DLQ para análise manual:
from celery import Celery
app = Celery('tasks', broker='redis://localhost:6379')
@app.task(bind=True, max_retries=3)
def send_email(self, email):
try:
provider.send(email)
except Exception as e:
if self.request.retries < self.max_retries:
# Retry
raise self.retry(countdown=2 ** self.request.retries)
else:
# 3 tentativas falharam, move pra DLQ
dlq.enqueue(email, error=str(e))
raise
# DLQ: fila separada pra jobs "mortos"
# Um admin verifica, investiga, e requeue manualmente se necessário
Ordenação vs Paralelismo
Ordenação garante: job A, depois B, depois C. Rápido? Não.
# ❌ Ordem estrita (sequencial)
queue = Queue(serializer='sequential') # A + B + C = 30 segundos
# ✓ Paralelo (rápido, mas sem ordem)
queue = Queue(serializer='parallel') # A, B, C simultâneos = 10 segundos
Trade-off: se precisa ordem (processamento de mudanças), use sequential. Se precisa velocidade (envio de emails), use paralelo.
Hybrid: ordenar por chave (partition key):
# Jobs do mesmo usuário: sequencial
# Jobs de usuários diferentes: paralelo
def enqueue_job(user_id, task):
queue.enqueue(task, partition_key=user_id)
# Redis roteriza:
# user_1 -> worker_1 (sequencial)
# user_2 -> worker_2 (paralelo)
# user_3 -> worker_3 (paralelo)
Monitorar profundidade da fila e lag
import prometheus_client
# Métrica: tamanho da fila
queue_size = prometheus_client.Gauge(
'job_queue_size',
'Número de jobs aguardando'
)
# Métrica: lag (atraso)
queue_lag = prometheus_client.Gauge(
'job_queue_lag_seconds',
'Tempo de espera no job mais antigo'
)
def monitor():
while True:
queue_size.set(redis.llen('celery')) # tamanho
# Job mais antigo na fila
oldest_job = redis.lindex('celery', 0)
if oldest_job:
enqueued_time = oldest_job['timestamp']
lag = (now - enqueued_time).total_seconds()
queue_lag.set(lag)
time.sleep(10) # monitora a cada 10 segundos
Alarmes:
- Queue size > 10.000: adiciona mais workers
- Lag > 60 segundos: sistema tá atrasado, escala
Aplicação em IA: inferência longa fora da requisição
Modelo de IA demora 30 segundos. Coloca fora da requisição:
# ANTES: requisição bloqueia
@app.post("/analyze-text")
def analyze(text):
# Inference demora 30 segundos
result = model.predict(text) # BLOQUEIA
return {"result": result}
# DEPOIS: fila
@app.post("/analyze-text")
def analyze(text):
job_id = uuid.uuid4()
queue.enqueue(run_inference, job_id, text)
return {"job_id": job_id, "status": "processing"}
@app.get("/result/{job_id}")
def get_result(job_id):
result = cache.get(f"inference:{job_id}")
if result:
return {"status": "done", "result": json.loads(result)}
else:
return {"status": "processing"}
def run_inference(job_id, text):
result = model.predict(text)
cache.setex(f"inference:{job_id}", 3600, json.dumps(result))
Cliente faz poll:
// Front-end
const jobId = await submitAnalysis(text); // responde em 100ms
while (true) {
const result = await getResult(jobId);
if (result.status === 'done') {
console.log(result.result);
break;
}
await sleep(1000); // tenta a cada 1 segundo
}
Armadilhas comuns
1. Fila sem limite (memory leak):
redis-cli
> LLEN celery
45000 # 45 mil jobs aguardando, Redis tá com 5 GB!
Configure limite:
# docker-compose.yml
redis:
image: redis:7
command: redis-server --maxmemory 2gb --maxmemory-policy allkeys-lru
2. Sem timeout na tarefa:
# ❌ Job pode ficar "preso" pra sempre
@app.task
def slow_task():
result = query_external_api() # pode falhar, timeout, etc
return result
# ✓ Com timeout
@app.task(time_limit=30) # falha após 30 segundos
def slow_task():
result = query_external_api()
return result
3. Retry infinito:
# ❌ Roda forever se erro persistente
@app.task(bind=True)
def always_fails(self):
raise Exception()
# self.retry() infinitamente
# ✓ Com limite
@app.task(bind=True, max_retries=5)
def always_fails(self):
raise Exception()
Quando NÃO usar filas
Tarefa rápida: se demora < 50 ms, fila adiciona overhead (serialization, rede). Roda na requisição mesmo.
Sem infraestrutura: fila precisa broker (RabbitMQ, Redis, AWS SQS). Se tá começando, talvez seja overengineering.
Dependência sequencial forte: se A depende de B depende de C, orquestração com fila fica complexa. Use job scheduler (DAG) em vez disso.
Próximos passos
Monitore filas em produção com Observabilidade em aplicações com LLM — inclui métricas de fila.
Deploy com infraestrutura: Terraform na prática: sua primeira infra versionada — provisiona RabbitMQ ou SQS.
Teste: cria Redis local, coloca fila em docker-compose.yml, escreve task que dorme 5 segundos, faz enqueue de 100 jobs, vê workers processarem em paralelo. Timing real é muito pedagógico.