Pular para o conteúdo
iauaiiauai — portal de tecnologia, IA e Cloud
EngenhariaIntermediário

Filas e workers: processando em background

Produtor/consumidor, idempotência obrigatória, retry backoff+jitter, DLQ. Ordenação vs paralelismo, lag. Código IA, Celery e Bull, garantias de entrega.

Por Equipe iauai · 4 de agosto de 2026 · 13 min de leitura

Nesta página

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.

Aprenda jogando

Pipeline em Pânico

Roteie registros sujos pelas estações certas antes que a esteira te engula.

Jogar Pipeline em Pânico

Continue lendo