0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

【LLM×非同期処理】深夜のAPIレートリミット(429 Too Many Requests)をCelery×Redisのトークンバケットで美しく制御する

0
Posted at

【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の定番非同期タスクキュー CeleryRedis を組み合わせ、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など)での実装例に興味がある方は、コメント欄までどうぞ。

0
0
0

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?