Tarefas em Segundo Plano e Celery — Filas de Tarefas Assíncronas
BackgroundTasks é como um sistema de senhas de restaurante—você recebe um número (resposta da API) imediatamente após o pedido, e é notificado quando a comida está pronta; Celery é como uma cozinha central—múltiplos chefes preparam diferentes pedidos simultaneamente, e o status de cada pedido pode ser rastreado, com retentativas em caso de falha.
1. O Que Você Vai Aprender
BackgroundTasksintegrado do FastAPI: Uma solução rápida para cenários simples- Arquitetura Celery: Worker / Broker (Redis) / Backend / Monitoramento Flower
- Integrando Celery com FastAPI: Definições de Tarefas, Gatilhos e Endpoints de Consulta de Status
- Retentativa de Tarefas e Tratamento de Erros:
@app.task(retry=3, acks_late=True) - Cenário Alice: Extração em Lote de Preços—A API retorna imediatamente um ID de tarefa, e o Celery Worker processa assincronamente a coleta de preços de milhões de produtos
2. A História Real da Alice
(1) Dor: Tarefas demoradas bloqueiam respostas da API
O PriceTracker da Alice precisa rastrear os preços de milhões de produtos em lote, e uma única extração leva 30 minutos. Se processado de forma síncrona, as requisições da API sofreriam timeout, e o frontend do Bob teria que esperar 30 minutos para receber uma resposta—resultando em uma péssima experiência do usuário. Enquanto isso, o BackgroundTasks do FastAPI só pode rodar dentro do processo atual, e reiniciar um worker causaria a perda da tarefa.
(2) Solução com Filas Distribuídas Celery
O Celery coloca tarefas demoradas em uma fila de mensagens (Redis Broker), onde processos Worker as consomem de forma assíncrona, e a API retorna imediatamente o ID da tarefa. As tarefas são automaticamente retentadas em caso de falha, os Workers podem ser escalados horizontalmente, e reinícios de processos não afetam as tarefas na fila.
from celery import Celery
celery_app = Celery("pricetracker", broker="redis://localhost:6379/0")
@celery_app.task(bind=True, max_retries=3)
def scrape_prices(self, product_ids: list[int]):
# Extração assíncrona de preços - roda no Celery Worker
...
(3) Resultado
A extração de um milhão de produtos passou de "bloquear a API por 30 minutos" para "retornar um ID de tarefa em 1 segundo + processamento em segundo plano." O número de workers pode ser escalado de 1 para 10, reduzindo o tempo de extração de 30 minutos para 3 minutos. As tarefas retentam automaticamente três vezes em caso de falha, aumentando a taxa de sucesso de 95% para 99,9%.
3. Soluções Leves para Tarefas em Segundo Plano
(1) Cenários Aplicáveis
(1) ▶ Exemplo: Enviando Notificações com BackgroundTasks
from fastapi import FastAPI, BackgroundTasks
from pydantic import BaseModel
app = FastAPI()
class PriceAlertRequest(BaseModel):
product_id: int
target_price: float
email: str
def send_price_alert_email(email: str, product_id: int, price: float):
# Simular envio de email (NÃO use await aqui)
print(f"Enviando alerta para {email}: Produto {product_id} atingiu ${price}")
@app.post("/alerts")
async def create_alert(alert: PriceAlertRequest, bg: BackgroundTasks):
# Adicionar tarefa para rodar após a resposta ser enviada
bg.add_task(send_price_alert_email, alert.email, alert.product_id, alert.target_price)
return {"message": "Alerta criado", "product_id": alert.product_id}
Saída:
# Função definida com sucesso
(2) Árvore de Decisão: BackgroundTasks vs. Celery
flowchart TD
Start{Precisa de tarefa em segundo plano?} --> Time{Leva > 1 min?}
Time -->|Não| Simple[Use BackgroundTasks]
Time -->|Sim| Retry{Precisa de retentativa/resiliência?}
Retry -->|Não| Simple
Retry -->|Sim| Scale{Precisa de escala horizontal?}
Scale -->|Não| Simple
Scale -->|Sim| Celery[Use Celery]
Simple -->|Prós| P1[Simples, sem infra]
Simple -->|Contras| C1[Sem retentativa, sem escala, perdido ao reiniciar]
Celery -->|Prós| P2[Retentativa, escala, persistência, monitoramento]
Celery -->|Contras| C2[Infra Redis + Worker]
| Dimensão | Tarefas em Segundo Plano | Celery |
|---|---|---|
| Complexidade | Zero configuração | Requer Redis + Worker |
| Persistência | Memória do Processo | Persistência Redis |
| Retentativa | Nenhuma | Mecanismo de retentativa integrado |
| Escalabilidade | Processo único | Escala horizontal baseada em Workers |
| Monitoramento | Nenhum | Flower Dashboard |
| Caso de Uso | Tarefas leves < 1 segundo | Tarefas demoradas > 1 minuto |
4. Arquitetura Celery
(1) Visão Geral da Arquitetura
flowchart TD
API[FastAPI App] -->|Enfileirar Tarefa| Broker[(Redis Broker)]
Broker -->|Consumir Tarefa| Worker1[Celery Worker 1]
Broker -->|Consumir Tarefa| Worker2[Celery Worker 2]
Broker -->|Consumir Tarefa| WorkerN[Celery Worker N]
Worker1 -->|Armazenar Resultado| Backend[(Redis Backend)]
Worker2 -->|Armazenar Resultado| Backend
WorkerN -->|Armazenar Resultado| Backend
API -->|Consultar Status| Backend
Flower[Flower Monitor] -->|Observar| Broker
Flower -->|Observar| Backend
| Componente | Função | Recomendação |
|---|---|---|
| Broker | Fila de Mensagens de Tarefas | Redis |
| Backend | Armazenamento de Resultados | Redis |
| Worker | Processo de Execução de Tarefas | celery -A app worker |
| Flower | Dashboard de Monitoramento | celery -A app flower |
(1) ▶ Exemplo: Configuração do Celery
# app/core/celery_app.py
from celery import Celery
celery_app = Celery(
"pricetracker",
broker="redis://localhost:6379/0",
backend="redis://localhost:6379/1",
)
celery_app.conf.update(
task_serializer="json",
accept_content=["json"],
result_serializer="json",
timezone="UTC",
enable_utc=True,
task_track_started=True,
task_acks_late=True, # Confirmar após execução, não antes
worker_prefetch_multiplier=4,
result_expires=3600, # Resultados expiram após 1 hora
)
Saída:
# Execução Bem-sucedida
5. Definição e Gatilho de Tarefas Celery
(1) Definição de Tarefa
(1) ▶ Exemplo: Tarefa de Extração de Preços
# app/tasks/price_scraping.py
from app.core.celery_app import celery_app
import asyncio
from sqlalchemy import select
@celery_app.task(bind=True, max_retries=3, default_retry_delay=60)
def scrape_product_prices(self, product_ids: list[int]):
"""Extrair preços atuais dos produtos dados."""
try:
for pid in product_ids:
# Simular extração (em produção: requisições HTTP para fontes de preços)
price = fetch_price_from_source(pid)
# Armazenar no banco de dados
save_price_record(pid, price)
return {"scraped": len(product_ids), "status": "success"}
except Exception as exc:
# Retentativa com backoff exponencial
raise self.retry(exc=exc, countdown=60 * (2 ** self.request.retries))
@celery_app.task(bind=True)
def bulk_scrape_all(self, total_products: int = 1000000, batch_size: int = 1000):
"""Extrair todos os produtos em lotes - padrão chord."""
batches = [
list(range(i, min(i + batch_size, total_products)))
for i in range(0, total_products, batch_size)
]
# Distribuir para tarefas de extração individuais
for batch in batches:
scrape_product_prices.delay(batch)
return {"total_batches": len(batches), "status": "started"}
Saída:
# Função definida com sucesso
(2) Transições de Status da Tarefa
stateDiagram-v2
[*] --> PENDING: Tarefa criada
PENDING --> STARTED: Worker pega a tarefa
STARTED --> PROGRESS: Em execução (opcional)
PROGRESS --> SUCCESS: Concluída
PROGRESS --> FAILURE: Erro ocorreu
STARTED --> FAILURE: Erro ocorreu
FAILURE --> RETRY: max_retries não atingido
RETRY --> PENDING: Re-enfileirada
FAILURE --> [*]: max_retries atingido
SUCCESS --> [*]
(2) ▶ Exemplo: Endpoint FastAPI Dispara uma Tarefa Celery
from fastapi import FastAPI, Depends
from app.core.celery_app import celery_app
from app.tasks.price_scraping import scrape_product_prices, bulk_scrape_all
app = FastAPI()
@app.post("/api/v1/scrape/prices")
async def trigger_scrape(
product_ids: list[int],
user=Depends(require_subscription("pro")),
):
# Disparar tarefa Celery - retorna ID da tarefa imediatamente
task = scrape_product_prices.delay(product_ids)
return {"task_id": task.id, "status": "pending"}
@app.post("/api/v1/scrape/bulk")
async def trigger_bulk_scrape(
user=Depends(require_subscription("enterprise")),
):
task = bulk_scrape_all.delay(total_products=1000000)
return {"task_id": task.id, "status": "pending"}
Saída:
# Função definida com sucesso
(3) ▶ Exemplo: Endpoint de Consulta de Status da Tarefa
from celery.result import AsyncResult
@app.get("/api/v1/tasks/{task_id}")
async def get_task_status(task_id: str):
result = AsyncResult(task_id, app=celery_app)
response = {
"task_id": task_id,
"status": result.status,
}
if result.ready():
if result.successful():
response["result"] = result.result
else:
response["error"] = str(result.result)
elif result.state == "PROGRESS":
response["progress"] = result.info
return response
Saída:
# Função definida com sucesso
6. Executando e Monitorando Celery Workers
(1) ▶ Exemplo: Iniciando Worker e Flower
# Iniciar Celery Worker
celery -A app.core.celery_app worker --loglevel=info --concurrency=4
# Iniciar dashboard de monitoramento Flower
celery -A app.core.celery_app flower --port=5555
# Visite http://localhost:5555 para o dashboard de monitoramento
Saída:
# Comando executado com sucesso
(2) ▶ Exemplo: Relatório de Progresso da Tarefa
from celery import current_task
@celery_app.task(bind=True)
def scrape_with_progress(self, product_ids: list[int]):
total = len(product_ids)
for i, pid in enumerate(product_ids):
# Processar cada produto
price = fetch_price_from_source(pid)
save_price_record(pid, price)
# Relatar progresso
self.update_state(
state="PROGRESS",
meta={"current": i + 1, "total": total, "percent": (i + 1) / total * 100},
)
return {"scraped": total}
Saída:
# Função definida com sucesso
❓ Perguntas Frequentes
task_acks_late=True?acks_late=True faz com que o worker confirme a tarefa apenas após a execução ser concluída; se travar, a tarefa será re-executada por outro worker.async/await em tarefas Celery?asyncio.run(). Alternativamente, use os modelos de concorrência eventlet ou gevent no Worker.📖 Resumo
- BackgroundTasks são adequadas para tarefas leves que levam menos de 1 segundo (envio de emails, escrita de logs); não exigem configuração, mas não suportam retentativas ou persistência.
- Celery é uma fila de tarefas para produção: Broker (fila Redis) + Worker (execução) + Backend (armazenamento de resultados)
@app.task(bind=True, max_retries=3)Definir tarefas retentáveis;self.retry()Disparar retentativa- FastAPI dispara uma tarefa via
.delay()e consulta o status viaAsyncResult; a API retorna imediatamente o ID da tarefa - Flower fornece um dashboard de monitoramento visual que permite visualizar o status dos workers, progresso das tarefas e taxas de sucesso em tempo real
📝 Exercícios
- Problema Básico (Dificuldade ⭐): Use
BackgroundTaskspara implementar um endpoint que envia um email de notificação em segundo plano (simulando saída de log) após um produto ser criado, e a API retorna imediatamente o resultado da criação. Dica:bg: BackgroundTasks+bg.add_task(fn, args) - Exercício Avançado (Dificuldade ⭐⭐): Configure o Celery (Redis Broker + Backend), defina a tarefa
scrape_prices(max_retries=3), faça um endpoint FastAPI disparar a tarefa e retornar o task_id, e outro endpoint para consultar o status da tarefa. Dica:celery_app.delay()+AsyncResult(task_id) - Desafio (Dificuldade: ⭐⭐⭐): Implemente uma tarefa de extração em lote com relatório de progresso—
scrape_with_progressrelata a porcentagem de progresso paraself.update_state(state="PROGRESS"), o frontend consulta/tasks/{task_id}para exibir uma barra de progresso, e usuários Enterprise podem disparar extração em nível de milhão. Dica: Useself.update_state()+result.infopara obter o progresso
---|



