Celery で Python の非同期タスク処理 — キュー、ワーカー、運用チェックリスト
メール送信、レポート生成、サムネイル変換、外部 API の大量呼び出し。こうした処理を Web リクエストの中でやるとレスポンスは遅くなり、タイムアウトが起き、失敗してもリトライする手段がありません。答えは昔から同じです。リクエストではジョブをキューに登録するだけにして、実際の実行は別プロセスに任せます。 Python でこのパターンの事実上の標準が Celery です。この記事では Celery の構造と基本的な使い方、そして運用で必ず引っかかるポイントを整理します。
構造 — ブローカー、ワーカー、結果バックエンド #
Celery は 3 つの部品で動きます。
| 部品 | 役割 | 代表的な選択肢 |
|---|---|---|
| ブローカー(broker) | タスクメッセージを入れるキュー | Redis、RabbitMQ |
| ワーカー(worker) | キューから取り出して実行するプロセス | celery worker コマンドで起動 |
| 結果バックエンド(result backend) | 実行結果・状態の保存(任意) | Redis、DB |
Web アプリが task.delay() を呼ぶとブローカーにメッセージが入り、ワーカーが取り出して実行します。Web プロセスとワーカープロセスは完全に分離されているので、それぞれ独立にスケールできます。ブローカーの選択はシンプルに整理できます。すでに Redis を使っているなら Redis で始め、メッセージの消失が絶対に許されない、あるいは複雑なルーティングが必要なら RabbitMQ を検討します。
最小構成 #
uv add "celery[redis]"# tasks.py
from celery import Celery
app = Celery(
"myapp",
broker="redis://localhost:6379/0",
backend="redis://localhost:6379/1",
)
@app.task
def send_welcome_email(user_id: int) -> str:
# 実際にはここでメール送信
return f"user {user_id} に送信完了"# ワーカーの起動
celery -A tasks worker --loglevel=info# Web 側のコード: 登録だけしてすぐ返す
result = send_welcome_email.delay(42)
print(result.id) # タスク ID
print(result.get(timeout=10)) # 結果が必要なら待機(Web リクエスト内では非推奨)delay() は即座に返ります。Web リクエストの中で result.get() を待つとタスクキューを使う意味がなくなるので、結果が必要な場合はタスク ID を返してクライアントが状態照会 API をポーリングする構成が一般的です。
リトライ — 失敗は基本前提です #
バックグラウンドジョブの相手はほとんどが外部の世界(メールサーバー、外部 API)なので、一時的な失敗は日常です。リトライはデコレーターのオプションで宣言します。
@app.task(
autoretry_for=(ConnectionError, TimeoutError), # この例外なら自動リトライ
max_retries=5,
retry_backoff=True, # 指数バックオフ: 1 秒、2 秒、4 秒...
retry_backoff_max=600, # バックオフの上限 10 分
retry_jitter=True, # ランダムな遅延を混ぜて同時リトライの殺到を防ぐ
)
def call_external_api(payload: dict) -> dict:
...- リトライ対象の例外を明示することが重要です。すべての例外をリトライすると、コードのバグ(KeyError など)まで 5 回繰り返し実行されます。
- 指数バックオフとジッターは、相手のサービスが落ちているときにリトライが殺到して二重に痛めつけるのを防ぐ基本装置です。
冪等性 — 最低 1 回実行を前提に書きます #
Celery の配送保証は基本的に最低 1 回(at-least-once)です。ワーカーが実行中に死ねば同じメッセージが別のワーカーで再実行されることがあり、リトライまで重なれば、同じタスクが 2 回走る状況はいつか必ず来ます。だからタスクは 2 回実行されても結果が同じになるように(冪等に) 書くのが原則です。
- 「ポイントを 1000 点追加」ではなく「注文 X に対するポイント付与を記録(すでにあれば無視)」として設計します。ユニーク制約や処理履歴テーブルが道具になります。
- 決済や発送のように冪等にしづらい処理は、外部サービスの冪等性キー(idempotency key)を活用します。
運用設定 — 事故の前にオンにしておくもの #
app.conf.update(
task_acks_late=True, # 実行完了後にキューから削除(ワーカー死亡時に再配送)
worker_prefetch_multiplier=1, # 長いタスクが混ざるなら先取りを減らす
task_time_limit=600, # 10 分を超えたら強制終了
task_soft_time_limit=540, # 9 分で例外を発生させて後片付けの機会を与える
)- acks_late: デフォルトは「受け取った時点で確認済み」なので、ワーカーが実行中に死ぬとそのジョブは消えます。オンにすると完了後の確認に変わる代わりに、上で述べた重複実行の可能性が生まれます。冪等性とセットの設定です。
- タイムリミット: 制限のないタスクがハングすると、ワーカーのスロットがひとつ永久に失われます。ソフトリミットで後片付けし、ハードリミットが最後の防衛線です。
- 監視: Flower を立てると、キューの長さ、タスクの成功・失敗、ワーカーの状態を Web UI で見られます。最低でもキューの長さにはアラートを掛けるべきです。キューが伸び続けるのは、ワーカーが処理量に追いついていないという最も早いシグナルです。
Celery まで必要ないケース #
Celery は強力ですが、ブローカーの運用、ワーカーのデプロイ、監視というインフラコストが付いてきます。もっと軽い選択肢が合う場合も多いです。
- FastAPI BackgroundTasks: レスポンス後に同じプロセスで実行される、いちばん軽い方式です。リトライも永続性もないので、失敗しても構わない副次処理(ログ蓄積、キャッシュ更新)までにとどめます。モダン Python 実践 #5 で扱いました。
- RQ、arq: Redis 専用のシンプルなタスクキューです。RQ は同期、arq は asyncio ベースです。設定が Celery よりはるかに少なく、タスクの種類が数個程度のプロジェクトなら十分です。
- Celery が合うケース: タスクの種類が多く、リトライ・スケジューリング(celery beat)・ルーティングといった機能が本当に必要で、運用を担えるチームがあるときです。
判断基準をひとつに絞るとこうなります。「このジョブが消えたら事故か?」 事故ならブローカーの耐久性とリトライを備えた Celery(または RabbitMQ の組み合わせ)にし、そうでなければもっと軽い道具で始めるほうが得策です。
まとめ #
- タスクキューの構造はブローカー(キュー)、ワーカー(実行)、結果バックエンド(状態)の 3 部品で、Web とワーカーは独立にスケールします。
delay()は登録だけして即座に返ります。Web リクエストの中で結果を待つ構成は避けます。- リトライは対象の例外を明示し、指数バックオフ + ジッターをオンにします。配送保証は最低 1 回なので、タスクは冪等に設計します。
- acks_late、タイムリミット、キュー長のアラートは、事故の前にオンにしておく基本設定です。
- 失敗しても構わない処理は BackgroundTasks、シンプルなキューは RQ・arq、機能と規模が必要になったら Celery です。