AIモデルとプラットフォーム
Pythonで非同期LLM API呼び出し:包括的なガイド
開発者およびデータサイエンティストとして、APIを介してこれらの強力なモデルとやり取りする必要が出てくることがよくあります。ただし、アプリケーションが複雑化し、スケールアップするにつれて、効率的でパフォーマンスの高いAPI呼び出しが必要になることが重要です。これが非同期プログラミングが輝く場所です。非同期プログラミングにより、LLM APIで作業するときに、最大のスループットと最小の待ち時間を実現できます。
この包括的なガイドでは、Pythonで非同期LLM API呼び出しについて探究します。非同期プログラミングの基礎から、複雑なワークフローを処理するための高度なテクニックまで、すべてをカバーします。この記事を読み終えるまでに、非同期プログラミングを使用してLLMパワードアプリケーションを強化する方法についてのsolidな理解を得ることができます。
非同期LLM API呼び出しについての詳細に踏み込む前に、非同期プログラミングの概念についてのsolidな基礎を確立しましょう。
非同期プログラミングにより、複数の操作をブロックせずに同時に実行できます。Pythonでは、主にasyncioモジュールを使用して、コルーチン、イベントループ、フューチャーを使用して同時実行コードを記述します。
重要な概念:
- コルーチン: async defで定義される関数で、停止して再開できます。
- イベントループ: 非同期タスクを管理して実行するための中心的な実行メカニズムです。
- 待ち可能なオブジェクト: awaitキーワードで使用できるオブジェクト(コルーチン、タスク、フューチャー)です。
これらの概念を示すために、以下は単純な例です:
import asyncio
<p>async def greet(name):
await asyncio.sleep(1) # I/O操作をシミュレート
print(f"こんにちは、{name}!")</p>
<p>async def main():
await asyncio.gather(
greet("Alice"),
greet("Bob"),
greet("Charlie")
)</p>
asyncio.run(main())
この例では、非同期関数greetを定義し、asyncio.sleep()を使用してI/O操作をシミュレートします。関数mainでは、asyncio.gather()を使用して複数の挨拶を同時に実行します。待ち時間があるにもかかわらず、すべての挨拶は約1秒後に印刷され、非同期実行の力が示されます。
非同期LLM API呼び出しの必要性
LLM APIを使用する場合、シーケンシャルまたは並列で複数のAPI呼び出しを行うシナリオに遭遇することがよくあります。従来の同期コードは、特にネットワークリクエストなどの高待ち時間の操作を扱う場合に、重大なパフォーマンスのボトルネックにつながる可能性があります。非同期アプローチを使用すると、複数のAPI呼び出しを同時に開始し、全体の実行時間を大幅に削減できます。
100個の異なる記事の要約をLLM APIを使用して生成するシナリオを考えてみましょう。同期アプローチでは、各API呼び出しは応答を待ってブロックされ、すべてのリクエストを完了するのに数分かかる可能性があります。一方、非同期アプローチを使用すると、複数のAPI呼び出しを同時に開始し、全体の実行時間を大幅に削減できます。
環境の設定
非同期LLM API呼び出しを開始するには、必要なライブラリを使用してPython環境を設定する必要があります。必要なものは次のとおりです:
- Python 3.7以上(ネイティブのasyncioサポートのため)
- aiohttp: 非同期HTTPクライアントライブラリ
- openai: OpenAIの公式Pythonクライアント(OpenAIのGPTモデルを使用する場合)
- langchain: LLMを使用したアプリケーションを構築するためのフレームワーク(オプションですが、複雑なワークフローには推奨)
これらの依存関係をpipを使用してインストールできます:
<p>pip install aiohttp openai langchain <div class="relative flex flex-col rounded-lg">
asyncioとaiohttpを使用した基本的な非同期LLM API呼び出し
OpenAIのGPT-3.5 APIを使用した単純な非同期LLM API呼び出しの例を見てみましょう:
import asyncio
import aiohttp
from openai import AsyncOpenAI
<p>async def generate_text(prompt, client):
response = await client.chat.completions.create(
model="gpt-3.5-turbo",
messages=[{"role": "user", "content": prompt}]
)
return response.choices[0].message.content</p>
<p>async def main():
prompts = [
"簡単な用語で量子コンピューティングを説明してください.",
"人工知能についての俳句を書いてください.",
"光合成のプロセスを説明してください."
]</p>
<p>async with AsyncOpenAI() as client:
tasks = [generate_text(prompt, client) for prompt in prompts]
results = await asyncio.gather(*tasks)</p>
<p>for prompt, result in zip(prompts, results):
print(f"プロンプト: {prompt}\nレスポンス: {result}\n")</p>
asyncio.run(main())
このアプローチにより、複数のリクエストをLLM APIに同時に送信し、全体の処理時間を大幅に削減できます。
高度なテクニック:バッチ処理と同時実行制御
実際のアプリケーションでは、より洗練されたアプローチが必要になることがよくあります。バッチ処理と同時実行制御の2つの重要なテクニックを探究してみましょう:
バッチ処理:大量のプロンプトを扱う場合、個別のリクエストではなくバッチで送信する方が効率的です。これにより、複数のAPI呼び出しのオーバーヘッドが削減され、パフォーマンスが向上します。
import asyncio
from openai import AsyncOpenAI
<p>async def process_batch(batch, client):
responses = await asyncio.gather(*[
client.chat.completions.create(
model="gpt-3.5-turbo",
messages=[{"role": "user", "content": prompt}]
) for prompt in batch
])
return [response.choices[0].message.content for response in responses]</p>
<p>async def main():
prompts = [f"番号{i}についての事実を教えてください" for i in range(100)]
batch_size = 10</p>
<p>async with AsyncOpenAI() as client:
results = []
for i in range(0, len(prompts), batch_size):
batch = prompts[i:i+batch_size]
batch_results = await process_batch(batch, client)
results.extend(batch_results)</p>
<p>for prompt, result in zip(prompts, results):
print(f"プロンプト: {prompt}\nレスポンス: {result}\n")</p>
asyncio.run(main())
同時実行制御:非同期プログラミングにより同時実行が可能ですが、APIサーバーを圧倒したり、レート制限を超過したりしないように同時実行レベルを制御することが重要です。asyncio.Semaphoreを使用して同時実行を制御できます。
import asyncio
from openai import AsyncOpenAI
<p>async def generate_text(prompt, client, semaphore):
async with semaphore:
response = await client.chat.completions.create(
model="gpt-3.5-turbo",
messages=[{"role": "user", "content": prompt}]
)
return response.choices[0].message.content</p>
<p>async def main():
prompts = [f"番号{i}についての事実を教えてください" for i in range(100)]
max_concurrent_requests = 5
semaphore = asyncio.Semaphore(max_concurrent_requests)</p>
<p>async with AsyncOpenAI() as client:
tasks = [generate_text(prompt, client, semaphore) for prompt in prompts]
results = await asyncio.gather(*tasks)</p>
<p>for prompt, result in zip(prompts, results):
print(f"プロンプト: {prompt}\nレスポンス: {result}\n")</p>
asyncio.run(main())
この例では、セマフォを使用して同時実行を制御し、APIサーバーを圧倒しないようにします。
非同期LLM呼び出しのエラーハンドリングとリトライ
外部APIを使用する場合、堅牢なエラーハンドリングとリトライメカニズムの実装が重要です。コードを強化して、一般的なエラーを処理し、指数関数的バックオフを使用してリトライを実装しましょう:
import asyncio
import random
from openai import AsyncOpenAI
from tenacity import retry, stop_after_attempt, wait_exponential
class APIError(Exception):
pass
<p>@retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=4, max=10))
async def generate_text_with_retry(prompt, client):
try:
response = await client.chat.completions.create(
model="gpt-3.5-turbo",
messages=[{"role": "user", "content": prompt}]
)
return response.choices[0].message.content
except Exception as e:
print(f"エラーが発生しました: {e}")
raise APIError("テキストの生成に失敗しました")</p>
<p>async def process_prompt(prompt, client, semaphore):
async with semaphore:
try:
result = await generate_text_with_retry(prompt, client)
return prompt, result
except APIError:
return prompt, "複数の試行後にレスポンスの生成に失敗しました"</p>
<p>async def main():
prompts = [f"番号{i}についての事実を教えてください" for i in range(20)]
max_concurrent_requests = 5
semaphore = asyncio.Semaphore(max_concurrent_requests)</p>
<p>async with AsyncOpenAI() as client:
tasks = [process_prompt(prompt, client, semaphore) for prompt in prompts]
results = await asyncio.gather(*tasks)</p>
<p>for prompt, result in results:
print(f"プロンプト: {prompt}\nレスポンス: {result}\n")</p>
asyncio.run(main())
この強化されたバージョンには、次のものが含まれています:
- API関連のエラー用のカスタム
APIError例外 - 指数関数的バックオフを使用してリトライを実装する
@retryデコレーターで装飾されたgenerate_text_with_retry関数 process_prompt関数内のエラーハンドリング
パフォーマンスの最適化:ストリーミングレスポンス
長いコンテンツを生成する場合、ストリーミングレスポンスはアプリケーションのパフォーマンスを大幅に改善できます。完全なレスポンスを待つ代わりに、利用可能になるにつれてテキストのチャンクを処理して表示できます。
import asyncio
from openai import AsyncOpenAI
<p>async def stream_text(prompt, client):
stream = await client.chat.completions.create(
model="gpt-3.5-turbo",
messages=[{"role": "user", "content": prompt}],
stream=True
)</p>
<p>full_response = ""
async for chunk in stream:
if chunk.choices[0].delta.content is not None:
content = chunk.choices[0].delta.content
full_response += content
print(content, end='', flush=True)</p>
<p>print("\n")
return full_response</p>
<p>async def main():
prompt = "時間旅行した科学者についての短い話を書いてください"</p>
<p>async with AsyncOpenAI() as client:
result = await stream_text(prompt, client)</p>
<p>print(f"完全なレスポンス:\n{result}")</p>
asyncio.run(main())
この例では、APIからレスポンスをストリーミングし、到着するたびにチャンクを印刷します。これは、チャットアプリケーションやユーザーにリアルタイムのフィードバックを提供したいシナリオで特に役立ちます。
LangChainを使用した非同期ワークフローの構築
より複雑なLLMパワードアプリケーションの場合、LangChainフレームワークは、複数のLLM呼び出しを連結し、他のツールを統合するプロセスを簡素化するための高い抽象化を提供します。LangChainを使用した非同期ワークフローの例を見てみましょう:
この例は、LangChainを使用してストリーミングと非同期実行を使用することで、より複雑なワークフローを作成する方法を示しています。AsyncCallbackManagerとStreamingStdOutCallbackHandlerを使用すると、生成されたコンテンツをリアルタイムでストリーミングできます。
import asyncio
from langchain.llms import OpenAI
from langchain.prompts import PromptTemplate
from langchain.chains import LLMChain
from langchain.callbacks.manager import AsyncCallbackManager
from langchain.callbacks.streaming_stdout import StreamingStdOutCallbackHandler
<p>async def generate_story(topic):
llm = OpenAI(temperature=0.7, streaming=True, callback_manager=AsyncCallbackManager([StreamingStdOutCallbackHandler()]))
prompt = PromptTemplate(
input_variables=["topic"],
template="{topic}についての短い話を書いてください"
)
chain = LLMChain(llm=llm, prompt=prompt)
return await chain.arun(topic=topic)</p>
<p>async def main():
topics = ["魔法の森", "未来の都市", "水中文明"]
tasks = [generate_story(topic) for topic in topics]
stories = await asyncio.gather(*tasks)</p>
<p>for topic, story in zip(topics, stories):
print(f"トピック: {topic}\nストーリー: {story}\n{'='*50}\n")</p>
asyncio.run(main())
FastAPIを使用した非同期LLMアプリケーションの提供
非同期LLMアプリケーションをWebサービスとして提供するには、FastAPIは非同期操作をネイティブにサポートしているため、優れた選択肢です。テキスト生成用の単純なAPIエンドポイントを作成する例を見てみましょう:
from fastapi import FastAPI, BackgroundTasks
from pydantic import BaseModel
from openai import AsyncOpenAI
app = FastAPI()
client = AsyncOpenAI()
<p>class GenerationRequest(BaseModel):
prompt: str</p>
<p>class GenerationResponse(BaseModel):
generated_text: str</p>
<p>@app.post("/generate", response_model=GenerationResponse)
async def generate_text(request: GenerationRequest, background_tasks: BackgroundTasks):
response = await client.chat.completions.create(
model="gpt-3.5-turbo",
messages=[{"role": "user", "content": request.prompt}]
)
generated_text = response.choices[0].message.content</p>
<p># 一部のポスト処理をバックグラウンドでシミュレート
background_tasks.add_task(log_generation, request.prompt, generated_text)</p>
<p>return GenerationResponse(generated_text=generated_text)</p>
<p>async def log_generation(prompt: str, generated_text: str):
# ロギングまたは追加処理をシミュレート
await asyncio.sleep(2)
print(f"ログ: プロンプト '{prompt}' が生成されたテキストの長さ {len(generated_text)}")</p>
<p>if __name__ == "__main__":
import uvicorn
uvicorn.run(app, host="0.0.0.0", port=8000)
このFastAPIアプリケーションは、プロンプトを受け取り、テキストを生成して返します。また、バックグラウンドタスクを使用して追加の処理を実行する方法も示しています。
ベストプラクティスと一般的な落とし穴
非同期LLM APIを使用する場合、以下のベストプラクティスを念頭に置いてください:
- 接続プーリングを使用する: 複数のリクエストを行う場合、接続を再利用してオーバーヘッドを削減します。
- 適切なエラーハンドリングを実装する: ネットワークの問題、APIエラー、予期しないレスポンスに対処します。
- レート制限を尊重する: セマフォやその他の同時実行制御メカニズムを使用して、APIサーバーを圧倒しないようにします。
- 監視とログを実装する: パフォーマンスを追跡し、問題を特定するための包括的なログを保持します。
- 長いコンテンツの場合にストリーミングを使用する: ユーザーにリアルタイムのフィードバックを提供し、部分的な結果の処理を可能にします。












