Chuyển đến nội dung chính

レッスン 15: バックグラウンド タスク、セロリ、タスク キュー

FastAPI BackgroundTasks、Redis/RabbitMQ ブローカーを備えた Celery、スケジュールされたタスク用の Celery Beat。タスクチェーン、エラー処理、再試行戦略。非同期タスク キューの ARQ および SAQ の代替手段。

💻 プログラミング — レッスン 15 レッスン 15: バックグラウンド タスク、セロリとタスク キュー

Python FastAPI: 基本から高度まで

パート 4: 高度な機能

xdev.asia

1. FastAPI バックグラウンドタスク

FastAPI には、応答の送信後に実行される軽量タスク用の BackgroundTasks が組み込まれています。

from fastapi import BackgroundTasks, FastAPI

app = FastAPI()


# Background task function
async def send_email(to: str, subject: str, body: str):
    """Gửi email - chạy sau khi response trả về."""
    import asyncio
    await asyncio.sleep(2)  # Simulate sending email
    print(f"Email sent to {to}: {subject}")


async def write_log(message: str):
    """Ghi log - chạy background."""
    with open("app.log", "a") as f:
        f.write(f"{message}\n")


@app.post("/users/")
async def create_user(
    email: str,
    background_tasks: BackgroundTasks,
):
    # Tạo user... (logic chính)
    user = {"id": 1, "email": email}

    # Thêm background tasks - chạy SAU khi response trả về
    background_tasks.add_task(send_email, email, "Welcome!", "Thanks for joining")
    background_tasks.add_task(write_log, f"User created: {email}")

    return user  # Response trả về NGAY, không đợi email


# Background tasks trong Dependencies
async def log_request(background_tasks: BackgroundTasks):
    """Dependency tự thêm background task."""
    background_tasks.add_task(write_log, "Request processed")

BackgroundTasks をいつ使用するか?

  • ✅ 軽くて速いタスク (数秒): 電子メールの送信、ログ、キャッシュの更新
  • ❌ 重くて長いタスク (分/時間): ML トレーニング、ビデオ処理 → Celery を使用する
  • ❌ タスクの再試行、スケジュール設定が必要 → Celery を使用

2.セロリのセットアップ

uv add celery[redis] redis
# app/core/celery_app.py
from celery import Celery

from app.config import settings

celery_app = Celery(
    "fastapi_app",
    broker=settings.redis_url,           # Message broker
    backend=settings.redis_url,          # Result backend
    include=["app.tasks.email", "app.tasks.reports"],
)

# Celery configuration
celery_app.conf.update(
    task_serializer="json",
    accept_content=["json"],
    result_serializer="json",
    timezone="UTC",
    enable_utc=True,
    task_track_started=True,
    task_time_limit=300,           # Hard limit: 5 minutes
    task_soft_time_limit=240,      # Soft limit: 4 minutes
    worker_prefetch_multiplier=1,  # 1 task at a time per worker
    worker_max_tasks_per_child=100,  # Restart worker after 100 tasks
)

3. セロリのタスク

# app/tasks/email.py
from celery import shared_task

from app.core.celery_app import celery_app


@celery_app.task(
    bind=True,
    max_retries=3,
    default_retry_delay=60,  # Retry sau 60 giây
    acks_late=True,
)
def send_email_task(self, to: str, subject: str, body: str):
    """Celery task gửi email."""
    try:
        # Send email logic
        import smtplib
        print(f"Sending email to {to}: {subject}")
        # ... actual email sending code
        return {"status": "sent", "to": to}
    except Exception as exc:
        # Retry with exponential backoff
        raise self.retry(exc=exc, countdown=2 ** self.request.retries * 60)


@celery_app.task(bind=True)
def send_bulk_emails(self, recipients: list[dict]):
    """Gửi email hàng loạt."""
    results = []
    for i, recipient in enumerate(recipients):
        try:
            send_email_task.delay(
                recipient["email"],
                recipient["subject"],
                recipient["body"],
            )
            results.append({"email": recipient["email"], "status": "queued"})
        except Exception as e:
            results.append({"email": recipient["email"], "status": "failed", "error": str(e)})

        # Update progress
        self.update_state(
            state="PROGRESS",
            meta={"current": i + 1, "total": len(recipients)},
        )
    return results
# app/tasks/reports.py
from celery import shared_task
from datetime import datetime

from app.core.celery_app import celery_app


@celery_app.task(bind=True, time_limit=600)
def generate_report(self, report_type: str, params: dict):
    """Generate báo cáo - task nặng."""
    self.update_state(state="PROGRESS", meta={"step": "Collecting data..."})

    # Simulate heavy computation
    import time
    time.sleep(5)

    self.update_state(state="PROGRESS", meta={"step": "Processing..."})
    time.sleep(5)

    self.update_state(state="PROGRESS", meta={"step": "Generating PDF..."})
    time.sleep(3)

    return {
        "report_type": report_type,
        "file_url": f"/reports/{report_type}_{datetime.utcnow().strftime('%Y%m%d')}.pdf",
        "generated_at": datetime.utcnow().isoformat(),
    }

4. Celery と FastAPI を統合する

# app/api/v1/tasks.py
from fastapi import APIRouter, HTTPException
from celery.result import AsyncResult

from app.core.celery_app import celery_app
from app.tasks.email import send_email_task
from app.tasks.reports import generate_report

router = APIRouter(prefix="/tasks", tags=["Tasks"])


@router.post("/send-email")
async def queue_email(to: str, subject: str, body: str):
    """Queue email task."""
    task = send_email_task.delay(to, subject, body)
    return {"task_id": task.id, "status": "queued"}


@router.post("/reports/{report_type}")
async def queue_report(report_type: str):
    """Queue report generation."""
    task = generate_report.delay(report_type, {})
    return {"task_id": task.id, "status": "queued"}


@router.get("/status/{task_id}")
async def get_task_status(task_id: str):
    """Kiểm tra status của task."""
    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.status == "PROGRESS":
        response["progress"] = result.info

    return response


@router.delete("/cancel/{task_id}")
async def cancel_task(task_id: str):
    """Hủy task đang chờ."""
    celery_app.control.revoke(task_id, terminate=True)
    return {"task_id": task_id, "status": "cancelled"}

5. Celery Beat - スケジュールされたタスク

# app/core/celery_app.py - thêm beat schedule
from celery.schedules import crontab

celery_app.conf.beat_schedule = {
    # Chạy mỗi 5 phút
    "cleanup-expired-tokens": {
        "task": "app.tasks.maintenance.cleanup_expired_tokens",
        "schedule": 300.0,  # 5 minutes
    },
    # Chạy lúc 2:00 AM hàng ngày
    "daily-report": {
        "task": "app.tasks.reports.generate_daily_report",
        "schedule": crontab(hour=2, minute=0),
    },
    # Chạy mỗi thứ 2 lúc 9:00 AM
    "weekly-digest": {
        "task": "app.tasks.email.send_weekly_digest",
        "schedule": crontab(hour=9, minute=0, day_of_week=1),
    },
}

6. セロリを実行する

# Chạy Celery worker
celery -A app.core.celery_app worker --loglevel=info --concurrency=4

# Chạy Celery Beat (scheduler)
celery -A app.core.celery_app beat --loglevel=info

# Chạy cả worker + beat
celery -A app.core.celery_app worker --beat --loglevel=info

# Monitor với Flower
pip install flower
celery -A app.core.celery_app flower --port=5555
# Dashboard: http://localhost:5555

7. Celery 用の Docker Compose

# docker-compose.yml
services:
  app:
    build: .
    ports:
      - "8000:8000"
    depends_on:
      - redis
      - postgres

  celery-worker:
    build: .
    command: celery -A app.core.celery_app worker --loglevel=info
    depends_on:
      - redis
      - postgres

  celery-beat:
    build: .
    command: celery -A app.core.celery_app beat --loglevel=info
    depends_on:
      - redis

  redis:
    image: redis:7-alpine
    ports:
      - "6379:6379"

  postgres:
    image: postgres:16-alpine
    environment:
      POSTGRES_DB: mydb
      POSTGRES_USER: user
      POSTGRES_PASSWORD: password
    ports:
      - "5432:5432"

概要

この記事では次のことを学びました。

  • バックグラウンドタスク: 軽作業用に内蔵
  • セロリ: 重いタスクのための分散タスクキュー
  • タスクの再試行: 指数バックオフ再試行戦略
  • タスクの進行状況: 長時間実行されるタスクの進行状況を追跡する
  • セロリビート: スケジュールされたタスクと定期的なタスク
  • Docker Compose: FastAPI を使用して Celery を実行する

次の記事では、ファイルのアップロード、キャッシュ、非同期の詳細について学習します。