非同期エージェントフレームワーク:LangChain とその先
LangChain と LangGraph における ainvoke()、astream()、非同期チェーンを扱います。
「非同期エージェントフレームワーク:LangChain とその先」はCoddyKit上の無料AI Agentsレッスンです。 これはレッスン4/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはAI Agents学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 AI Agentsコースには全4レッスンが含まれています。
LangChainでの非同期実行
LangChainは、すべてのインターフェースに非同期版を提供しています。invoke()を持つすべてのコンポーネントにはainvoke()があり、stream()にはastream()があります。本番環境のエージェントには、非同期処理が推奨されます。
非同期LLM呼び出しのainvoke()
ainvoke()はinvoke()に相当する非同期版です。非同期関数内で使用して、ノンブロッキングなLLM呼び出しを行います。これにより、複数のエージェントやリクエストでイベントループを共有できます。
import asyncio
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage
llm = ChatOpenAI(model='gpt-4o-mini', api_key='sk-...')
async def async_agent_call(question: str) -> str:
# ainvoke: non-blocking, releases event loop while waiting for OpenAI
response = await llm.ainvoke([HumanMessage(content=question)])
return response.content
async def handle_multiple_users(questions: list) -> list:
# All three LLM calls run concurrently
results = await asyncio.gather(*[async_agent_call(q) for q in questions])
return results
questions = [
'What is Python?',
'What is TypeScript?',
'What is Rust?'
]
results = asyncio.run(handle_multiple_users(questions))
for q, a in zip(questions, results):
print(f'Q: {q[:30]}... A: {a[:50]}...')トークンストリーミングのastream()
astream()は、LLMから到着したトークンを順次生成します。これにより、完全な生成結果を待たずに、ユーザーへレスポンスをリアルタイムでストリーミングできます。
import asyncio
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage
llm = ChatOpenAI(model='gpt-4o-mini', api_key='sk-...')
async def stream_response(question: str):
print(f'Streaming answer to: {question}\n')
full_response = ''
async for chunk in llm.astream([HumanMessage(content=question)]):
token = chunk.content
if token:
print(token, end='', flush=True) # Print each token as it arrives
full_response += token
print() # New line after streaming
return full_response
async def main():
await stream_response('List 3 benefits of async programming in Python')
asyncio.run(main())AsyncCallbackHandler
LangChainのコールバックは、LLMの開始、LLMの終了、ツールの開始、チェーンエラーなど、特定のイベントで発生します。AsyncCallbackHandlerは、エージェントのループをブロックせずに、これらのイベントを非同期で処理します。
from langchain_core.callbacks import AsyncCallbackHandler
from typing import Any, Dict, List
import time
class LatencyCallbackHandler(AsyncCallbackHandler):
def __init__(self):
self.step_times = {}
self.step_counts = {}
async def on_llm_start(self, serialized: Dict, prompts: List[str], **kwargs):
run_id = str(kwargs.get('run_id', ''))
self.step_times[run_id] = time.perf_counter()
async def on_llm_end(self, response, **kwargs):
run_id = str(kwargs.get('run_id', ''))
if run_id in self.step_times:
elapsed_ms = (time.perf_counter() - self.step_times[run_id]) * 1000
print(f'LLM call completed in {elapsed_ms:.0f}ms')
async def on_tool_start(self, serialized: Dict, input_str: str, **kwargs):
tool_name = serialized.get('name', 'unknown')
print(f'Tool starting: {tool_name}')
async def on_tool_error(self, error: Exception, **kwargs):
print(f'Tool error: {error}')
handler = LatencyCallbackHandler()
print('Async callback handler created')
# Use: llm.ainvoke([...], config={'callbacks': [handler]})LangGraphの非同期ノード関数
LangGraphのノードは非同期関数にできます。ノードをasync defとして定義すると、LangGraphはグラフの実行中にそのノードをawaitします。本番環境のグラフには、このパターンが推奨されます。
import asyncio
from langgraph.graph import StateGraph, END
from typing import TypedDict, List
class AgentState(TypedDict):
question: str
entities: List[str]
context: str
answer: str
async def extract_entities_node(state: AgentState) -> AgentState:
await asyncio.sleep(0.1) # Simulate async NLP call
entities = state['question'].split()[:3] # Simplified
return {'entities': entities}
async def retrieve_context_node(state: AgentState) -> AgentState:
await asyncio.sleep(0.2) # Simulate async vector search
context = f'Context for entities: {state["entities"]}'
return {'context': context}
async def generate_answer_node(state: AgentState) -> AgentState:
await asyncio.sleep(0.3) # Simulate async LLM call
answer = f'Answer based on: {state["context"]}'
return {'answer': answer}
# Build async graph
graph = StateGraph(AgentState)
graph.add_node('extract', extract_entities_node)
graph.add_node('retrieve', retrieve_context_node)
graph.add_node('generate', generate_answer_node)
graph.set_entry_point('extract')
graph.add_edge('extract', 'retrieve')
graph.add_edge('retrieve', 'generate')
graph.add_edge('generate', END)
app = graph.compile()
print('Async LangGraph compiled')LangGraphからの非同期ストリーミング
LangGraphは、グラフの実行中に中間状態を非同期でストリーミングできます。完全な実行が終わるまで待つのではなく、各ノードの出力が完了するたびに確認するには、astream()を使用します。
import asyncio
async def stream_graph_execution(graph_app, initial_state: dict):
print('Graph execution streaming:')
async for step_output in graph_app.astream(initial_state):
for node_name, state_delta in step_output.items():
print(f' Node [{node_name}] completed:')
for key, value in state_delta.items():
print(f' {key}: {value}')
# Run the async graph
initial = {
'question': 'What is machine learning?',
'entities': [],
'context': '',
'answer': ''
}
asyncio.run(stream_graph_execution(app, initial))セマフォによるレート制限
OpenAIなどのLLM APIには、1分あたりのリクエスト数に関するレート制限があります。非同期セマフォを使用すると、多数のエージェントタスクを同時に実行する場合でも、レート制限を超えないようにできます。
import asyncio
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage
llm = ChatOpenAI(model='gpt-4o-mini', api_key='sk-...')
# Limit to 10 concurrent LLM calls
LLM_SEMAPHORE = asyncio.Semaphore(10)
async def rate_limited_llm_call(question: str) -> str:
async with LLM_SEMAPHORE:
response = await llm.ainvoke([HumanMessage(content=question)])
return response.content
async def process_large_batch(questions: list) -> list:
print(f'Processing {len(questions)} questions with max 10 concurrent LLM calls')
tasks = [rate_limited_llm_call(q) for q in questions]
results = await asyncio.gather(*tasks, return_exceptions=True)
successes = [r for r in results if not isinstance(r, Exception)]
failures = [r for r in results if isinstance(r, Exception)]
print(f'Success: {len(successes)}, Failed: {len(failures)}')
return results
# Process 50 questions with max 10 concurrent calls
questions = [f'Question {i}: What is concept number {i}?' for i in range(20)]
asyncio.run(process_large_batch(questions))非同期ツールの定義
LangChainでは、ツール関数を非同期にできます。非同期ツールはエージェントの実行中にawaitされるため、ツール自体でノンブロッキングなAPI呼び出しを実行できます。
import asyncio
import httpx
from langchain.tools import tool
@tool
async def async_web_search(query: str) -> str:
'''Search the web for information about the query.'''
async with httpx.AsyncClient() as client:
# Real implementation would use a search API
response = await client.get(
'https://api.search.example.com/search',
params={'q': query, 'api_key': 'your-key'},
timeout=10.0
)
response.raise_for_status()
results = response.json()
return '\n'.join([r['snippet'] for r in results.get('items', [])[:3]])
@tool
async def async_fetch_document(url: str) -> str:
'''Fetch and return the text content of a URL.'''
async with httpx.AsyncClient() as client:
response = await client.get(url, timeout=15.0)
return response.text[:3000] # Limit content size
print('Async tools defined')
print('Use with: agent.ainvoke({"input": "your question"})')OpenAIを直接使った非同期エージェント
LangChainを使わず、OpenAI SDKだけで完全に非同期なエージェントループを直接構築できます。これにより、最大限の制御と最小限のオーバーヘッドを実現できます。
import asyncio
import openai
import json
client = openai.AsyncOpenAI(api_key='sk-...')
TOOLS = [
{'type': 'function', 'function': {
'name': 'web_search',
'description': 'Search the web',
'parameters': {'type': 'object', 'properties': {'query': {'type': 'string'}}, 'required': ['query']}
}}
]
async def async_tool_call(tool_name: str, args: dict) -> str:
if tool_name == 'web_search':
await asyncio.sleep(0.3) # Simulate search
return f'Search results for: {args["query"]}'
return 'Unknown tool'
async def async_agent_loop(question: str, max_turns: int = 5) -> str:
messages = [{'role': 'user', 'content': question}]
for turn in range(max_turns):
response = await client.chat.completions.create(
model='gpt-4o-mini', messages=messages, tools=TOOLS
)
msg = response.choices[0].message
messages.append(msg)
if not msg.tool_calls:
return msg.content
# Execute tool calls in parallel
tool_results = await asyncio.gather(*[
async_tool_call(tc.function.name, json.loads(tc.function.arguments))
for tc in msg.tool_calls
])
for tc, result in zip(msg.tool_calls, tool_results):
messages.append({'role': 'tool', 'tool_call_id': tc.id, 'content': result})
return 'Max turns reached'
result = asyncio.run(async_agent_loop('What is the latest news on AI?'))
print(result)キャンセルとクリーンアップ
非同期タスクはキャンセルできます。エージェントの実行がキャンセルされた場合(ユーザーの要求やタイムアウトなど)にリソースが確実にクリーンアップされるよう、asyncio.CancelledErrorを適切に処理してください。
import asyncio
async def cancellable_agent(question: str):
try:
print('Agent starting')
await asyncio.sleep(0.5) # Step 1
print('Step 1 done')
await asyncio.sleep(0.5) # Step 2 - may be cancelled here
print('Step 2 done')
return 'Completed'
except asyncio.CancelledError:
print('Agent was cancelled - cleaning up')
# Clean up resources: close connections, log cancellation
raise # Always re-raise CancelledError
finally:
print('Cleanup always runs')
async def run_with_timeout(question: str, timeout: float):
task = asyncio.create_task(cancellable_agent(question))
try:
result = await asyncio.wait_for(task, timeout=timeout)
return result
except asyncio.TimeoutError:
print(f'Agent exceeded {timeout}s timeout')
task.cancel()
return None
# Run with 0.7s timeout (not enough for both steps)
result = asyncio.run(run_with_timeout('test', timeout=0.7))
print('Final result:', result)非同期エージェントコードのテスト
pytest-asyncioを使用して、非同期エージェント関数をテストします。テスト関数に@pytest.mark.asyncioを付けると、イベントループ内で実行できます。
import pytest
import asyncio
from unittest.mock import AsyncMock, patch
# Install: pip install pytest-asyncio
# pytest.ini: [pytest] asyncio_mode = auto
@pytest.mark.asyncio
async def test_async_agent_call():
with patch('openai.AsyncOpenAI') as mock_openai:
mock_client = AsyncMock()
mock_openai.return_value = mock_client
mock_response = AsyncMock()
mock_response.choices[0].message.content = 'Mocked answer'
mock_response.choices[0].message.tool_calls = None
mock_client.chat.completions.create.return_value = mock_response
# Test the async function
result = await async_agent_call('What is Python?')
assert isinstance(result, str)
print('Async test passed')
@pytest.mark.asyncio
async def test_parallel_execution():
start = asyncio.get_event_loop().time()
results = await asyncio.gather(
asyncio.sleep(0.1),
asyncio.sleep(0.1),
asyncio.sleep(0.1)
)
elapsed = asyncio.get_event_loop().time() - start
assert elapsed < 0.3, 'Should complete in parallel'
print(f'Parallel test passed: {elapsed:.2f}s')理解度チェック:非同期フレームワーク
非同期エージェントフレームワークについての理解度を確認します。
非同期フレームワークのまとめ
非同期版LangChainは、本番環境に対応した非同期エージェント向けに、ainvoke()、astream()、AsyncCallbackHandlerを提供します。LangGraphは非同期ノード関数をネイティブにサポートします。レート制限にはセマフォ、直接呼び出しにはAsyncOpenAIクライアント、テストにはpytest-asyncioを使用します。適切なキャンセル処理により、エージェントが中断された場合でもリソースを確実にクリーンアップできます。
よくある質問
「非同期エージェントフレームワーク:LangChain とその先」レッスンは無料ですか?
はい。「非同期エージェントフレームワーク:LangChain とその先」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、AI Agentsコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 AI Agentsコースには全4レッスンが含まれています。
「非同期エージェントフレームワーク:LangChain とその先」で何を学びますか?
LangChain と LangGraph における ainvoke()、astream()、非同期チェーンを扱います。 ブラウザで直接実行するハンズオンコードでAI Agentsを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。
AI Agentsを始めるのに経験は必要ですか?
事前経験は必要ありません。CoddyKitのAI Agentsは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン4/4です。
「非同期エージェントフレームワーク:LangChain とその先」レッスンにはどのくらい時間がかかりますか?
ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。
このAI Agentsレッスンでコードを書いて実行できますか?
はい。すべてのAI Agentsレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- エージェント開発者のための非同期 Python
- イベントキューとメッセージブローカー
- ノンブロッキングなツールの並列実行
- 非同期エージェントフレームワーク:LangChain とその先