サイトアイコン IT & ライフハックブログ|学びと実践のためのアイデア集

FastAPIの重い処理をCeleryとRedisへ分離する方法:再試行、定期実行、監視まで

green snake

Photo by Pixabay on Pexels.com

画像変換、帳票作成、外部APIとの連携などに時間がかかると、HTTPリクエストの中で処理を終える設計には限界があります。応答が遅くなるだけでなく、通信が切れたときの再実行や、処理状況を追う方法も決めなければなりません。

そこで、FastAPIで構築したWeb APIはジョブの受付に集中させ、実処理をCeleryワーカーへ渡します。Redisは、ジョブを受け渡すブローカーと、状態や結果を保存する結果バックエンドに利用します。本記事では、最小構成を動かすだけでなく、重複実行、再試行、監視まで含めて設計する手順を説明します。

構成要素と役割

最初に、同じRedisを使っていても、ブローカーと結果バックエンドは別の役割だと理解しておく必要があります。

構成要素 役割 停止したときの主な影響
FastAPI 入力を検証し、ジョブをキューへ登録して受付結果を返す 新しいジョブを受け付けられない
Redisブローカー APIからワーカーへジョブメッセージを受け渡す ジョブの投入や取得が止まる
Celeryワーカー キューからジョブを取り出して処理する 待機中のジョブが進まない
結果バックエンド ジョブの状態、進捗、戻り値を保存する 状態確認や結果取得ができない
Celery Beat 定期ジョブを決めた時刻にキューへ登録する 定期ジョブが新たに投入されない
Flower ワーカーやタスクを確認する運用画面 監視画面が使えないが、ワーカー自体は別に動作する

Redis URLの末尾を/1/2に分ける方法は、キー空間を整理するための論理分離です。同じRedisサーバーを使う限り、障害点やメモリ上限まで分離されるわけではありません。負荷、保持期間、可用性の要件が異なる場合は、Redisインスタンスやサービスそのものを分ける判断が必要です。

BackgroundTasksとCeleryの使い分け

FastAPIのBackgroundTasksは、レスポンスを返した後に同じアプリケーションプロセスで処理を続ける仕組みです。同じプロセスの情報を使う小さな後処理には簡潔ですが、別サーバーのワーカーへ仕事を配るジョブキューではありません。

判断項目 BackgroundTasksが合う場合 Celeryが合う場合
処理時間 短く、アプリの再起動時に失われても業務影響を管理できる 長時間になり、APIプロセスから切り離したい
再試行 独自の再試行制御が不要 一時障害を条件付きで再試行したい
実行場所 APIと同じプロセスでよい 別プロセスや別サーバーへ配りたい
運用 キューの監視が不要 待ち件数、失敗、実行時間を追跡したい
定期実行 別の仕組みで管理する Celery Beatなどと組み合わせたい

Celeryを導入すれば、すべてのジョブが一度だけ実行されるわけではありません。通信断、ワーカー停止、クライアントの再送によって、同じ業務処理が複数回試行される可能性を前提にします。そのため、ジョブキューの採用と冪等性の設計は一組です。

ジョブが完了するまでの流れ

  1. クライアントがFastAPIへ処理を依頼する。
  2. FastAPIが入力、権限、重複した依頼でないかを確認する。
  3. FastAPIがCeleryタスクをキューへ登録し、HTTP 202とジョブIDを返す。
  4. Celeryワーカーがジョブを取得し、処理を実行する。
  5. ワーカーが状態や結果を更新する。
  6. クライアントが認証済みの状態確認APIから進捗または結果を取得する。

HTTP 202は「処理を受け付けた」ことを示す応答です。業務処理の成功までは表しません。受付直後に成功扱いへ変えるのではなく、ジョブIDと状態確認方法を返す設計にします。

最小構成を実装する

必要なパッケージ

fastapi
uvicorn[standard]
celery[redis]
pydantic
pydantic-settings

設定を環境変数から読み込む

pydantic-settingsでは、設定項目の型を宣言し、環境変数や.envから値を読み込みます。現在の書き方ではSettingsConfigDictを使います。

# app/settings.py
from functools import lru_cache

from pydantic_settings import BaseSettings, SettingsConfigDict


class Settings(BaseSettings):
    app_name: str = "FastAPI Job API"
    celery_broker_url: str = "redis://redis:6379/1"
    celery_result_backend: str = "redis://redis:6379/2"

    model_config = SettingsConfigDict(
        env_file=".env",
        env_file_encoding="utf-8",
        extra="ignore",
    )


@lru_cache
def get_settings() -> Settings:
    return Settings()

.envを本番コンテナへそのまま組み込むのではなく、実行環境のシークレット管理方法に合わせて値を渡します。ブローカーや結果バックエンドの認証情報をログへ出さないことも必要です。

Celeryアプリとタスクを定義する

# app/tasks.py
from celery import Celery

from app.settings import get_settings

settings = get_settings()
celery_app = Celery(
    "job_app",
    broker=settings.celery_broker_url,
    backend=settings.celery_result_backend,
)
celery_app.conf.update(
    task_serializer="json",
    accept_content=["json"],
    result_serializer="json",
    task_track_started=True,
    timezone="UTC",
    enable_utc=True,
)


class RetryableJobError(Exception):
    """通信断など、再試行してよい一時障害を表す。"""


def build_report(payload: dict) -> dict:
    # 実際の帳票作成処理を実装する。
    return {"ok": True, "job": payload}


@celery_app.task(
    bind=True,
    autoretry_for=(RetryableJobError,),
    retry_backoff=True,
    retry_jitter=True,
    retry_kwargs={"max_retries": 5},
)
def generate_report(self, payload: dict) -> dict:
    return build_report(payload)

autoretry_for=(Exception,)のように対象を広げると、入力不備や実装ミスまで繰り返すおそれがあります。再試行してよい一時障害だけを専用の例外へ変換し、認証失敗、入力不備、恒久的な制約違反は失敗として終了させます。

FastAPIから投入して状態を確認する

# app/main.py
from celery.result import AsyncResult
from fastapi import FastAPI, status
from pydantic import BaseModel, Field

from app.settings import get_settings
from app.tasks import celery_app, generate_report

app = FastAPI(title=get_settings().app_name)


class ReportRequest(BaseModel):
    customer_id: int = Field(gt=0)
    report_type: str = Field(min_length=1, max_length=40)


@app.post("/jobs", status_code=status.HTTP_202_ACCEPTED)
def enqueue_job(payload: ReportRequest):
    task = generate_report.delay(payload.model_dump())
    return {"job_id": task.id, "state": "QUEUED"}


@app.get("/jobs/{job_id}")
def get_job(job_id: str):
    result = AsyncResult(job_id, app=celery_app)

    if result.state == "SUCCESS":
        return {"state": result.state, "result": result.result}
    if result.state == "FAILURE":
        return {"state": result.state, "error": "ジョブに失敗しました"}
    return {"state": result.state}

この状態確認APIは説明用の最小例です。本番では、ジョブIDだけで結果を取得させず、依頼者とジョブの対応をアプリケーションのデータベースへ保存して、認証と認可を確認します。CeleryのPENDINGは、待機中だけでなく、結果バックエンドがそのIDを知らない場合にも返り得るため、アプリ側のジョブ記録と照合する必要があります。

Composeで開発環境を起動する

DockerfileとComposeの役割を分けて考えると、API、ワーカー、スケジューラーを同じイメージから起動しやすくなります。現在のCompose仕様では、先頭のversion指定は不要です。

# compose.yaml
services:
  redis:
    image: redis:7-alpine
    command: ["redis-server", "--appendonly", "yes"]
    volumes:
      - redis-data:/data
    healthcheck:
      test: ["CMD", "redis-cli", "ping"]
      interval: 5s
      timeout: 3s
      retries: 10
    restart: unless-stopped

  app:
    build: .
    command: uvicorn app.main:app --host 0.0.0.0 --port 8000
    env_file: .env
    depends_on:
      redis:
        condition: service_healthy
    ports:
      - "8000:8000"
    restart: unless-stopped

  worker:
    build: .
    command: celery -A app.tasks.celery_app worker --loglevel=INFO
    env_file: .env
    depends_on:
      redis:
        condition: service_healthy
    restart: unless-stopped

  beat:
    build: .
    command: celery -A app.tasks.celery_app beat --loglevel=INFO
    env_file: .env
    depends_on:
      redis:
        condition: service_healthy
    restart: unless-stopped

volumes:
  redis-data:

この構成はローカル検証の出発点です。本番では、Redisの永続化、バックアップ、認証、通信経路、容量監視、復旧手順を要件に合わせて設計します。コンテナが再起動する設定だけでは、ジョブの保持や業務処理の整合性まで保証できません。

再試行と冪等性を別々に設計する

指数バックオフは、再試行の間隔を段階的に延ばす方法です。ジッタは、その間隔に揺らぎを加え、多数のワーカーが同時に再試行する集中を避けるために使います。どちらも一時障害への負荷を抑える仕組みであり、恒久的なエラーを成功へ変える仕組みではありません。

失敗の種類 扱い
一時障害 短い通信断、一時的な接続失敗、上流サービスの一時停止 回数と上限時間を決めて再試行する
恒久的な失敗 入力不備、権限不足、存在しない対象 再試行せず、修正可能な情報を運用側へ残す
原因不明 想定外の例外 失敗として記録し、調査後に再投入を判断する

冪等性は、同じ依頼が繰り返されても、請求、登録、通知などの業務上の副作用を重複させない性質です。APIで冪等キーを受け取り、依頼内容とジョブIDの組を一意に保存します。データベース更新には一意制約やアップサートを使い、外部サービスが冪等キーを受け付ける場合は同じキーを引き継ぎます。Webhookを扱う場合は、署名検証、再送対策、冪等性の設計も合わせて確認してください。

Redisロックは同時実行を抑える補助手段であり、完了済みの業務を記録する仕組みの代わりにはなりません。ロックを使うなら、処理時間を考慮した有効期限、所有者ごとに異なるトークン、所有者が一致するときだけ解放する処理が必要です。固定値を保存して最後に単純なDELを実行すると、期限切れ後に別のワーカーが取得したロックまで消す可能性があります。

定期実行と時刻指定の注意点

定期ジョブはCelery Beatのbeat_scheduleへ登録できます。Beatは同じスケジュールを二重投入しないよう、原則として一つのスケジューラーだけを動かします。ただし、前回の処理が終わる前に次の時刻を迎えることはあるため、タスク側の冪等性や排他制御は別に必要です。

from celery.schedules import crontab

celery_app.conf.beat_schedule = {
    "daily-report": {
        "task": "app.tasks.generate_report",
        "schedule": crontab(hour=6, minute=0),
        "args": ({"report_type": "daily"},),
    }
}

countdownetaは短い遅延には使えますが、指定時刻ちょうどの実行を保証するものではありません。待機中のタスクはワーカーのメモリに保持され、Redisブローカーではvisibility timeoutとの関係で再配送される場合があります。遠い将来の大量予約には使わず、永続化されたスケジュールや定期実行の仕組みを選びます。タイムゾーンは設定と運用時刻の基準を明示し、夏時間を含む地域時刻を扱う場合はテストします。

進捗、結果、データベースを管理する

bind=Trueで定義したタスク内からself.update_state(state="PROGRESS", meta={...})を使うと、結果バックエンドへ独自の進捗を保存できます。ただし、この呼び出しだけでブラウザーへ通知が届くわけではありません。ポーリングで取得するか、Celeryの状態を受け取ってWebSocketやServer-Sent Eventsへ渡す別の通知経路を用意します。

結果をいつまで保持するかも先に決めます。大きな成果物をRedisへ直接保存するのではなく、オブジェクトストレージやデータベースへ保存し、結果バックエンドには状態と参照先だけを置く方法があります。保存期間を過ぎた結果を消す処理と、アプリケーション側のジョブ履歴をどう残すかは別々に設計します。

ワーカーでは、APIリクエスト中に作ったデータベースセッションを受け渡しません。タスクごとにセッションを開き、成功時にコミットし、失敗時にロールバックして閉じます。

@celery_app.task
def save_report(job_id: int) -> None:
    db = SessionLocal()
    try:
        # 対象を取得し、結果を保存する。
        db.commit()
    except Exception:
        db.rollback()
        raise
    finally:
        db.close()

長い外部API呼び出しのあいだ、データベーストランザクションを開いたままにしないよう処理を分けます。ジョブ状態の更新と業務データの更新について、どこまでを一つのトランザクションにするかを決め、途中失敗から再開しても重複しない形にします。

監視とセキュリティ

FlowerはCeleryクラスターの状態を確認するための運用画面です。タスク引数や失敗情報が見える場合があり、タスクの停止などの操作機能も持つため、インターネットへ無認証で公開する構成は避けます。内部ネットワークに置き、認証を設定し、閲覧だけで足りる利用者には読み取り専用を検討します。

ダッシュボードを置くだけでは監視になりません。次の指標にしきい値と通知先を設定します。

APIではPydanticモデルで入力を検証し、タスクへ渡す値を最小限にします。巨大なファイル本体、認証情報、データベースセッションをメッセージへ入れず、アクセス制御された保存先の識別子を渡します。状態確認APIではジョブの所有者を検証し、失敗レスポンスには内部の例外文や接続先をそのまま含めません。

テストは業務処理とキュー連携を分ける

Celeryタスクの中に業務ロジックを詰め込むと、単体テストがキューの状態に依存します。帳票作成やデータ変換を通常の関数またはサービスへ分け、タスクは引数の受け取り、再試行、状態更新を担当させます。

同期的にタスク関数を呼ぶテストだけでは、実ワーカーとのシリアライズや配送の差を検出できません。速い単体テストと、対象を絞った統合テストを併用します。

本番導入前の確認項目

FastAPIを本番運用へ進める確認項目も使い、API本体の認証、ログ、終了処理、デプロイと合わせて点検すると抜けを減らせます。

ジョブ基盤は失敗後の動きまで決める

FastAPIから重い処理を外す目的は、応答を速く見せることだけではありません。受付と実行を分離し、一時障害を安全に再試行し、重複した依頼でも業務データを壊さず、運用担当者が滞留や失敗を把握できる状態を作ることです。

まずは一つのジョブを別ワーカーで実行し、202応答と状態確認を実装します。次に冪等性、失敗分類、監視を整え、実測に基づいてキュー分割やワーカー数を調整します。この順序なら、構成だけが複雑になり、障害時の判断材料がない状態を避けられます。

この記事に関連する株式会社greedenの取り組み

非同期ジョブ基盤は、失敗時の再実行や監視まで含めて設計してこそ運用できます。株式会社greedenは、要件整理からAPI開発、クラウド構築、保守改善まで、Webシステム開発を一気通貫で支援します。

モバイルバージョンを終了