画像変換、帳票作成、外部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を導入すれば、すべてのジョブが一度だけ実行されるわけではありません。通信断、ワーカー停止、クライアントの再送によって、同じ業務処理が複数回試行される可能性を前提にします。そのため、ジョブキューの採用と冪等性の設計は一組です。
ジョブが完了するまでの流れ
- クライアントがFastAPIへ処理を依頼する。
- FastAPIが入力、権限、重複した依頼でないかを確認する。
- FastAPIがCeleryタスクをキューへ登録し、HTTP 202とジョブIDを返す。
- Celeryワーカーがジョブを取得し、処理を実行する。
- ワーカーが状態や結果を更新する。
- クライアントが認証済みの状態確認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"},),
}
}
countdownとetaは短い遅延には使えますが、指定時刻ちょうどの実行を保証するものではありません。待機中のタスクはワーカーのメモリに保持され、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クラスターの状態を確認するための運用画面です。タスク引数や失敗情報が見える場合があり、タスクの停止などの操作機能も持つため、インターネットへ無認証で公開する構成は避けます。内部ネットワークに置き、認証を設定し、閲覧だけで足りる利用者には読み取り専用を検討します。
ダッシュボードを置くだけでは監視になりません。次の指標にしきい値と通知先を設定します。
- キューの待ち件数と最も古いジョブの待機時間
- 実行中、成功、失敗、再試行の件数
- タスクの実行時間とタイムアウト件数
- ワーカーの稼働数と応答
- Redisのメモリ使用量、接続失敗、永続化状態
APIではPydanticモデルで入力を検証し、タスクへ渡す値を最小限にします。巨大なファイル本体、認証情報、データベースセッションをメッセージへ入れず、アクセス制御された保存先の識別子を渡します。状態確認APIではジョブの所有者を検証し、失敗レスポンスには内部の例外文や接続先をそのまま含めません。
テストは業務処理とキュー連携を分ける
Celeryタスクの中に業務ロジックを詰め込むと、単体テストがキューの状態に依存します。帳票作成やデータ変換を通常の関数またはサービスへ分け、タスクは引数の受け取り、再試行、状態更新を担当させます。
- 単体テスト:業務関数へ入力を渡し、戻り値と副作用を検証する。
- タスク統合テスト:テスト用のRedisと実ワーカーを起動し、シリアライズ、配送、結果保存を確認する。
- APIテスト:202応答、入力エラー、認証、ジョブ所有者の照合を確認する。
- 障害テスト:一時障害の再試行、上限到達、ワーカー再起動、同じ依頼の重複投入を確認する。
同期的にタスク関数を呼ぶテストだけでは、実ワーカーとのシリアライズや配送の差を検出できません。速い単体テストと、対象を絞った統合テストを併用します。
本番導入前の確認項目
- 再試行する例外と、再試行しない例外を分けたか
- 同じ依頼や同じメッセージが複数回来ても副作用が重複しないか
- ジョブIDと依頼者を保存し、状態確認APIで認可しているか
- ブローカーと結果バックエンドの停止、容量超過、復旧手順を決めたか
- タスクの時間制限、ワーカーの並列度、キュー分割を実測で調整したか
- 定期ジョブが重なったときの動作を決めたか
- FlowerとRedisを無認証で外部公開していないか
- 待ち時間、失敗率、実行時間、ワーカー停止を通知できるか
FastAPIを本番運用へ進める確認項目も使い、API本体の認証、ログ、終了処理、デプロイと合わせて点検すると抜けを減らせます。
ジョブ基盤は失敗後の動きまで決める
FastAPIから重い処理を外す目的は、応答を速く見せることだけではありません。受付と実行を分離し、一時障害を安全に再試行し、重複した依頼でも業務データを壊さず、運用担当者が滞留や失敗を把握できる状態を作ることです。
まずは一つのジョブを別ワーカーで実行し、202応答と状態確認を実装します。次に冪等性、失敗分類、監視を整え、実測に基づいてキュー分割やワーカー数を調整します。この順序なら、構成だけが複雑になり、障害時の判断材料がない状態を避けられます。
この記事に関連する株式会社greedenの取り組み
非同期ジョブ基盤は、失敗時の再実行や監視まで含めて設計してこそ運用できます。株式会社greedenは、要件整理からAPI開発、クラウド構築、保守改善まで、Webシステム開発を一気通貫で支援します。
