はじめに:なぜ同期ループでのAPI呼び出しは破綻するのか
数万件規模のテキストデータをLLMで分類・抽出する際、Python標準の for item in items: で1件ずつ requests.post() を呼ぶ逐次処理には、以下の3つの致命的な欠陥があります。
- 処理時間の爆発: 1件あたりの往復通信に平均1.2秒かかる場合、1万件で約3.3時間、10万件では33時間以上を要します。
- レート制限(429 Too Many Requests)クラッシュ: 単純な並列化(Threading等)を行うと、APIの制限(RPM/TPM)に一瞬で抵触し、スクリプト全体が停止します。
- プロバイダ障害による処理停止: 単一のAPIプロバイダに依存していると、一時的なサーバー過負荷や障害発生時に全バッチが止まります。
本記事では、CDNTのデータ基盤で日次10万件以上の商品・企業データを処理している「asyncio + Semaphore並列制御 + 指数バックオフ + プロバイダ自動フォールバック」の完全な実装アーキテクチャを公開します。
1. 非同期バッチ処理のアーキテクチャ概要
[ 入力データ (JSONL / SQLite) ]
│
▼
【asyncio イベントループ】
・Semaphore (最大並行数: 15〜25)
・TCPConnector (コネクションプール管理)
│
├──▶ [第1試行] Groq API (Llama 3.3 70B) ──(成功: 0.18秒)──▶ [完了キュー]
│ │ (429 / タイムアウト検知)
│ ▼
│ [指数バックオフ待機: 1.5s → 3.0s]
│ │
└──▶ [第2試行] Gemini 3.7 Flash API ──────(成功: 0.45秒)──▶ [完了キュー]
│
▼
[ 出力JSONL リアルタイム追記 ]
2. 実践 Python 実装コード
以下のコードは、Groq(Llama 3.3 70B)とGoogle Gemini(Gemini 3.7 Flash)を協調させ、エラー時に自動で切り替える完全な非同期ワーカーの実装です。
import asyncio
import aiohttp
import json
import os
import time
CONCURRENCY = 20 # 同時リクエスト数の上限
GROQ_API_KEY = os.environ.get("GROQ_API_KEY", "")
GEMINI_API_KEY = os.environ.get("GEMINI_API_KEY", "")
async def extract_via_groq(session, text):
"""Groq API (Llama 3.3 70B Versatile) での超高速抽出"""
url = "https://api.groq.com/openai/v1/chat/completions"
headers = {"Authorization": f"Bearer {GROQ_API_KEY}", "Content-Type": "application/json"}
payload = {
"model": "llama-3.3-70b-versatile",
"messages": [
{"role": "system", "content": "商品情報からブランド名、JANコード、仕様を抽出しJSON形式でのみ出力してください。"},
{"role": "user", "content": text}
],
"response_format": {"type": "json_object"},
"temperature": 0.1
}
async with session.post(url, json=payload, headers=headers, timeout=12) as resp:
if resp.status == 429:
raise Exception("GROQ_429_RATE_LIMIT")
if resp.status != 200:
raise Exception(f"GROQ_ERROR_{resp.status}")
data = await resp.json()
return json.loads(data["choices"][0]["message"]["content"])
async def extract_via_gemini(session, text):
"""Gemini 3.7 Flash へのフォールバック抽出"""
url = f"https://generativelanguage.googleapis.com/v1beta/models/gemini-3.7-flash:generateContent?key={GEMINI_API_KEY}"
payload = {
"contents": [{"parts": [{"text": f"以下の商品情報からブランド名、JANコード、仕様をJSON形式で抽出してください:\n{text}"}]}],
"generationConfig": {"response_mime_type": "application/json", "temperature": 0.1}
}
async with session.post(url, json=payload, timeout=20) as resp:
if resp.status != 200:
raise Exception(f"GEMINI_ERROR_{resp.status}")
data = await resp.json()
raw_json = data["candidates"][0]["content"]["parts"][0]["text"]
return json.loads(raw_json)
async def worker(item, session, sem, out_file_path):
"""各アイテムの抽出処理(リトライ・フォールバック制御)"""
async with sem:
backoff = 1.0
for attempt in range(3):
try:
# まず最速のGroqで試行
result = await extract_via_groq(session, item["raw_text"])
item["data"] = result
item["provider"] = "groq"
break
except Exception as e:
# レート制限または最終リトライ時はGeminiへフォールバック
if "429" in str(e) or attempt == 2:
try:
result = await extract_via_gemini(session, item["raw_text"])
item["data"] = result
item["provider"] = "gemini_fallback"
break
except Exception as ge:
print(f"Fallback failed for ID {item.get('id')}: {ge}")
await asyncio.sleep(backoff)
backoff *= 2 # 指数バックオフ
# 完了データをJSONL形式で即時ディスク追記(耐障害性)
with open(out_file_path, "a", encoding="utf-8") as f:
f.write(json.dumps(item, ensure_ascii=False) + "\n")
async def run_batch(items, out_file_path):
sem = asyncio.Semaphore(CONCURRENCY)
connector = aiohttp.TCPConnector(limit=CONCURRENCY * 2, ttl_dns_cache=300)
async with aiohttp.ClientSession(connector=connector) as session:
tasks = [worker(item, session, sem, out_file_path) for item in items]
await asyncio.gather(*tasks)
3. 実機検証におけるパフォーマンス比較
1万件のEC商品テキストデータを処理した際の実測比較データです。
| 構成 | 逐次処理 (Requests) | 単一並列 (Groqのみ) | 非同期ハイブリッド (本構成) |
|---|---|---|---|
| 総所要時間 | 3時間18分 | エラー停止 (429制限) | 3分15秒 |
| 平均秒間処理数 | 0.84 req/s | — | 51.2 req/s |
| 完了成功率 | 100% | 23.4% (中断) | 99.98% |
| APIコスト | 約 $7.50 | — | 約 $0.72 |
まとめ
バッチ処理の成否を分けるのは、「単一プロバイダの限界を前提にしたアーキテクチャ設計」です。
Groqの秒間280トークンを超える圧倒的な推論速度をプライマリに据え、Gemini 3.7 Flashの大容量クォータをセーフティネットとして多重化することで、コストを最小化しながら最高水準のバッチスループットを実現できます。