404 Not Found

404 Not Found


nginx

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


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.

PYTHON
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

PYTHON
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:

TEXT
# Função definida com sucesso

(2) Árvore de Decisão: BackgroundTasks vs. Celery

100%
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

100%
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

PYTHON
# 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:

TEXT
# 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

PYTHON
# 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:

TEXT
# Função definida com sucesso

(2) Transições de Status da Tarefa

100%
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

PYTHON
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:

TEXT
# Função definida com sucesso

(3) ▶ Exemplo: Endpoint de Consulta de Status da Tarefa

PYTHON
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:

TEXT
# Função definida com sucesso

6. Executando e Monitorando Celery Workers

(1) ▶ Exemplo: Iniciando Worker e Flower

BASH
# 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:

TEXT
# Comando executado com sucesso

(2) ▶ Exemplo: Relatório de Progresso da Tarefa

PYTHON
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:

TEXT
# Função definida com sucesso

❓ Perguntas Frequentes

P Quando as BackgroundTasks são executadas?
R Elas são executadas após a resposta ser enviada ao cliente. Se uma tarefa lançar uma exceção, isso não afeta a resposta já enviada, mas será registrado no log.
P O Celery Worker e o FastAPI devem rodar no mesmo processo?
R Não. O Worker é um processo separado que é implantado e escalado independentemente. O FastAPI é apenas responsável por disparar tarefas, enquanto o Worker é responsável por executá-las.
P Qual é a diferença entre usar Redis como Broker e como Backend?
R O Broker é uma fila de tarefas (tarefas aguardando execução), enquanto o Backend é um armazenamento de resultados (resultados de tarefas concluídas). Você pode usar bancos de dados diferentes na mesma instância Redis (por exemplo, DB 0 e DB 1).
P Qual é o propósito de task_acks_late=True?
R Por padrão, um worker confirma uma tarefa assim que a recebe; se travar durante a execução, a tarefa é perdida. Configurar 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.
P Posso usar async/await em tarefas Celery?
R Tarefas Celery são funções síncronas. Para chamar código assíncrono, envolva-o em asyncio.run(). Alternativamente, use os modelos de concorrência eventlet ou gevent no Worker.
P Como lidar com o rastreamento de progresso para tarefas envolvendo milhões de itens?
R Use o padrão "chord": Divida a tarefa em 1.000 subtarefas (cada uma processando 1.000 itens), e quando as subtarefas forem concluídas, o callback "chord" agrega os resultados. O frontend consulta o status das tarefas "chord".

📖 Resumo


📝 Exercícios

  1. Problema Básico (Dificuldade ⭐): Use BackgroundTasks para 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)
  2. 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)
  3. Desafio (Dificuldade: ⭐⭐⭐): Implemente uma tarefa de extração em lote com relatório de progresso—scrape_with_progress relata a porcentagem de progresso para self.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: Use self.update_state() + result.info para obter o progresso

---|

Web-Tutorial.com

Equipe Técnica Web-Tutorial

Uma plataforma de tutoriais mantida por diversos desenvolvedores. Cada tutorial é escrito e revisado por profissionais da área correspondente. Trabalhamos para manter nosso conteúdo preciso e confiável — se encontrar algum problema, avise-nos.

100%