【LLM×非同期処理】深夜のAPIレートリミット(429 Too Many Requests)をCelery×Redisのトークンバケットで美しく制御する
はじめに
社内データを一括でベクトル化(Embedding)したり、数万件の商品データにLLMで自動タグ付けするような「大量バッチ処理」を実装したとき、高確率で我々を絶望させるのが 429 Too Many Requests (Rate Limit Exceeded) です。
OpenAI等のAPIには、1分あたりのリクエスト数(RPM)だけでなく、**1分あたりの消費トークン数(TPM)**という厳しい制限があります。
単純なループ処理や、安易な asyncio.gather で並列実行すると、一瞬でこの制限に引っかかり、タスクが途中で壊滅します。
本記事では、Pythonの定番非同期タスクキュー Celery と Redis を組み合わせ、APIの枠(RPM/TPM)を限界まで使い切りつつ、絶対に429エラーを出さない「トークンバケット型」の流量制御(Rate Limiting)を実装する方法を解説します。
1. なぜ「単なる time.sleep()」や「指数リトライ」ではダメなのか?
- time.sleepの限界: LLMのレスポンス時間は毎回変わるため、一律で3秒待つような実装にすると、遅すぎるか、結局上限に引っかかるかの二択になります。
- 指数バックオフ(リトライ)の罠: エラーを検知してから引く(待つ)戦略ですが、大量のタスクが同時にリトライに回ると、今度は「リトライの波」が重なってさらに大きな429のスパイクを叩き出す原因になります(群集事故)。
理想の解決策:
リクエストを投げる「前」に、Redisで一元管理された「トークン(許可証)」をタスクが取得し、許可が出た分だけAPIを叩くトークンバケットアルゴリズムをタスクキューに組み込みます。
2. システム構成
フロー図 (Mermaid)
3. 【実装】CeleryとRedisによるカスタムRate Limiter
今回は、Redisの機能を使って「1分間に最大10,000トークン、最大50リクエスト」という枠を、複数のWorker間で共有・制御するコードを構築します。
依存関係
poetry add celery redis openai tiktoken
実装コード (tasks.py)
import os
import time
from celery import Celery
from redis import Redis
from openai import OpenAI
import tiktoken
# RedisとCeleryの初期化
REDIS_URL = os.getenv("REDIS_URL", "redis://localhost:6379/0")
app = Celery("llm_tasks", broker=REDIS_URL, backend=REDIS_URL)
redis_client = Redis.from_url(REDIS_URL)
openai_client = OpenAI(api_key=os.getenv("OPENAI_API_KEY"))
encoder = tiktoken.encoding_for_model("gpt-4o-mini")
# APIの制限定義 (例: 本番環境のプランに合わせて調整)
MAX_RPM = 50
MAX_TPM = 10000
def check_and_consume_limits(estimated_tokens: int) -> bool:
"""
Redisのカウンターを使って、現在のRPM/TPMの枠に収まるかチェックする。
収まる場合は枠を消費してTrueを返す。
"""
current_minute = int(time.time() // 60)
rpm_key = f"llm_limit:rpm:{current_minute}"
tpm_key = f"llm_limit:tpm:{current_minute}"
# パイプラインでアトミックに処理
pipe = redis_client.pipeline()
pipe.get(rpm_key)
pipe.get(tpm_key)
current_rpm, current_tpm = pipe.execute()
current_rpm = int(current_rpm) if current_rpm else 0
current_tpm = int(current_tpm) if current_tpm else 0
# 制限プレチェック
if current_rpm + 1 > MAX_RPM or current_tpm + estimated_tokens > MAX_TPM:
return False # 枠がいっぱい
# 枠を消費
pipe = redis_client.pipeline()
pipe.incr(rpm_key)
pipe.expire(rpm_key, 120) # 2分後に自動消去
pipe.incrby(tpm_key, estimated_tokens)
pipe.expire(tpm_key, 120)
pipe.execute()
return True
@app.task(bind=True, max_retries=None)
def process_llm_heavy_task(self, text_content: str):
"""
大量に並列実行される重いLLMタスク
"""
# 1. 投げる前に、プロンプトのトークン数を概算(少し多めに見積もるのがコツ)
prompt_tokens = len(encoder.encode(text_content))
estimated_total_tokens = prompt_tokens + 500 # 出力分のバッファ
# 2. 枠があるか確認
if not check_and_consume_limits(estimated_total_tokens):
# 枠がない場合は、10秒後にこのタスクをリトライ(キューに優しく戻す)
# これによりWorkerのスレッドをロック(sleep)せずに解放できる
print(f"【流量制御】制限に達したため、タスクを退避します(想定トークン: {estimated_total_tokens})")
raise self.retry(countdown=10)
# 3. 枠が確保できた安全な状態でAPIを叩く
try:
response = openai_client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": text_content}]
)
return response.choices[0].message.content
except Exception as e:
# 万が一、他の要因でエラーが出た場合のセーフティネット
print(f"APIエラー発生: {e}")
raise self.retry(countdown=30)
4. この構成が優れている理由(実務上のメリット)
① Workerをノンブロッキングに保てる
time.sleep() でWorkerを停止させると、その間Workerのプロセスが丸々1つ無駄になります。この実装では、枠がないタスクは速やかにキューの最後尾(指定秒後)にリトライとして回り、Workerは「今すぐ処理できる別のタスク」を瞬時に引き受けられるため、インフラのCPU効率が最大化します。
② 分散システム(マルチサーバー)で完全同期する
Redisを中央カウンターにしているため、Workerが1台だろうが、オートスケーリングで10台に増えようが、システム全体での「合計TPM/RPM」を確実にコントロールできます。
まとめ
LLMのバッチ処理を安定して回すためには、APIの「外側」にあるインフラ層での流量制御が不可欠です。
429エラーを「受けてから慌てる」のではなく、Redisで「投げる前に止める」設計に切り替えるだけで、深夜のバッチ処理が驚くほど静かに、かつ最速で終わるようになります。
大量データ処理に悩むバックエンドエンジニアの参考になれば幸いです!
「このインフラ構成、知りたかった!」という方は、LGTM と ストック をお願いします!
Celery以外のキュー(SQSやCloud Tasksなど)での実装例に興味がある方は、コメント欄までどうぞ。