0Pricing
AI Engineering Academy · レッスン

ストリーミングレスポンスのツール呼び出しを処理する

関数呼び出しの引数がトークン単位で届くストリーミングレスポンスを解析し、JSONの断片をバッファリングして、呼び出しが完了したときだけツールを実行します。

「ストリーミングレスポンスのツール呼び出しを処理する」はCoddyKit上の無料AI Engineering Academyレッスンです。 これはレッスン4/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはAI Engineering Academy学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 AI Engineering Academyコースには全4レッスンが含まれています。

ストリーム内のツール呼び出しは異なる形で届く

LLMが関数の呼び出しを決定すると、レスポンスの構造が変わります。content文字列の代わりに、差分にtool_calls配列が含まれます。ただし、ストリーミングレスポンスでは、関数呼び出しの引数が不完全なJSON文字列としてトークン単位で届きます。1つのチャンクで完全なJSONオブジェクトを受け取るわけではありません。ツール呼び出しを解析して実行するには、これらの断片をバッファーに蓄積し、完全なJSONを再構成する必要があります。

# In a non-streaming response, tool call is complete:
# choice.message.tool_calls[0].function.arguments = '{"city": "Paris"}'

# In a streaming response, arguments arrive in pieces:
# chunk 1: delta.tool_calls[0].function.arguments = '{'
# chunk 2: delta.tool_calls[0].function.arguments = '"city"'
# chunk 3: delta.tool_calls[0].function.arguments = ': "'
# chunk 4: delta.tool_calls[0].function.arguments = 'Paris'
# chunk 5: delta.tool_calls[0].function.arguments = '"}'
# You must concatenate these before JSON.parse can work

ストリーム内のツール呼び出しを検出する

各チャンクの finish_reason を確認して、ツール呼び出しが発生するタイミングを判断します。finish_reason が 'tool_calls' の場合、モデルは関数を呼び出すことを決定しており、ストリームは終了します。finish_reason が 'stop' の場合、モデルは通常のテキスト応答を生成しています。ストリーミング中は、ツール呼び出しの引数の断片を識別するために、chunk.choices[0].delta.tool_calls が非 None かどうかを確認します。

async def detect_stream_type(messages, tools):
    stream = await async_client.chat.completions.create(
        model='gpt-4o-mini',
        messages=messages,
        tools=tools,
        stream=True,
    )
    response_type = 'text'
    async for chunk in stream:
        choice = chunk.choices[0]
        if choice.delta.tool_calls:  # tool call fragment arriving
            response_type = 'tool_call'
        if choice.finish_reason == 'tool_calls':
            print('Model wants to call a function')
        elif choice.finish_reason == 'stop':
            print('Normal text response')
    return response_type

ツール呼び出しの引数をバッファリングする

各チャンクから届く引数の断片を蓄積するには、ツール呼び出しのインデックスをキーとする辞書を使用します。このインデックスは、ストリーミング中のどのツール呼び出しかを識別します。1つの応答でモデルが複数の関数を呼び出すこともあります。tool_calls デルタが null でないチャンクごとに、引数の断片をツール呼び出しのインデックスに対応するバッファエントリへ追加します。

from collections import defaultdict

async def collect_streamed_tool_calls(messages, tools):
    stream = await async_client.chat.completions.create(
        model='gpt-4o-mini',
        messages=messages,
        tools=tools,
        stream=True,
    )

    tool_call_buffers = defaultdict(lambda: {'name': '', 'id': '', 'arguments': ''})
    text_buffer = ''

    async for chunk in stream:
        delta = chunk.choices[0].delta

        if delta.content:  # text content
            text_buffer += delta.content

        if delta.tool_calls:
            for tc in delta.tool_calls:
                idx = tc.index
                if tc.id:
                    tool_call_buffers[idx]['id'] = tc.id
                if tc.function.name:
                    tool_call_buffers[idx]['name'] += tc.function.name
                if tc.function.arguments:
                    tool_call_buffers[idx]['arguments'] += tc.function.arguments

    return text_buffer, dict(tool_call_buffers)

ツール呼び出しの解析と実行

ストリームが終了して引数の文字列が完全に揃ったら、それぞれを json.loads で解析し、適切な Python 関数に振り分けます。ツール呼び出しは要求された順序で実行するか、互いに独立している場合は並列で実行します。その後、結果をツール応答メッセージとして整形し、後続の API 呼び出しで送信します。

import json

# Example tool registry
tools_registry = {
    'get_weather': lambda city, unit='celsius': {'temp': 22, 'desc': 'sunny', 'city': city},
    'search_docs': lambda query, top_k=3: [{'title': 'Doc 1', 'snippet': 'Relevant info...'}],
}

def execute_tool_calls(tool_call_buffers: dict) -> list[dict]:
    tool_messages = []
    for idx in sorted(tool_call_buffers.keys()):
        tc = tool_call_buffers[idx]
        func_name = tc['name']
        args = json.loads(tc['arguments'])

        if func_name in tools_registry:
            result = tools_registry[func_name](**args)
        else:
            result = {'error': f'Unknown function: {func_name}'}

        tool_messages.append({
            'role': 'tool',
            'tool_call_id': tc['id'],
            'content': json.dumps(result),
        })
    return tool_messages

ストリーミングによるツール呼び出しの全体ループ

ストリーミングでのツール呼び出しを完全に処理するには、複数ターンの会話ループが必要です。最初のリクエストでツール呼び出しが返された場合は、それを実行して結果をメッセージ履歴に追加します。2回目のリクエストで最終的なテキスト回答が返されます。モデルが先行するツール呼び出しの結果に基づいて追加のツール呼び出しを選択した場合、このループが複数回繰り返されることもあります。

async def streaming_agent_loop(initial_messages, tools):
    messages = list(initial_messages)
    max_iterations = 5

    for iteration in range(max_iterations):
        text, tool_calls = await collect_streamed_tool_calls(messages, tools)

        if tool_calls:
            # Append assistant message with tool calls
            assistant_msg = {
                'role': 'assistant',
                'content': text or None,
                'tool_calls': [
                    {'id': tc['id'], 'type': 'function',
                     'function': {'name': tc['name'], 'arguments': tc['arguments']}}
                    for tc in tool_calls.values()
                ]
            }
            messages.append(assistant_msg)

            # Execute tools and append results
            tool_results = execute_tool_calls(tool_calls)
            messages.extend(tool_results)
        else:
            # No more tool calls — final text response
            print('Final answer:', text)
            return text

    return 'Max iterations reached'

ツール呼び出しをバッファリングしながらテキストをストリーミングする

実際のアプリケーションでは、ツール呼び出しの引数を同時にバッファリングしながら、テキストをすぐにクライアントへストリーミングしたい場合が多くあります。そのためには、content を含むチャンクはすぐにストリーミングし、tool_calls を含むチャンクは後で実行するためにバッファリングする必要があります。ツール呼び出しが完了して実行できる状態になるのは、ストリームが終了した後です。

async def stream_with_tools(messages, tools):
    stream = await async_client.chat.completions.create(
        model='gpt-4o-mini', messages=messages, tools=tools, stream=True
    )
    tool_buffers = defaultdict(lambda: {'name': '', 'id': '', 'arguments': ''})
    text_parts = []

    async for chunk in stream:
        delta = chunk.choices[0].delta
        finish = chunk.choices[0].finish_reason

        if delta.content:
            text_parts.append(delta.content)
            yield ('text', delta.content)   # stream to client immediately

        if delta.tool_calls:
            for tc in delta.tool_calls:
                if tc.id:    tool_buffers[tc.index]['id'] = tc.id
                if tc.function.name: tool_buffers[tc.index]['name'] += tc.function.name
                if tc.function.arguments: tool_buffers[tc.index]['arguments'] += tc.function.arguments

        if finish == 'tool_calls':
            yield ('tool_calls', dict(tool_buffers))  # signal tool execution needed

ツール呼び出しの並列実行

モデルが複数のツール呼び出しを同時に返す場合(並列関数呼び出しと呼ばれる機能)は、順番に実行するのではなく asyncio.gather で並列実行します。順次実行すると不要な遅延が発生します。たとえばモデルが天気 API とデータベース検索を同時に呼び出す場合、一方の完了を待ってからもう一方を開始する理由はありません。

import asyncio

async def execute_tool_calls_parallel(tool_call_buffers: dict) -> list[dict]:
    async def execute_one(idx, tc):
        func_name = tc['name']
        args = json.loads(tc['arguments'])
        if func_name in async_tools_registry:
            result = await async_tools_registry[func_name](**args)
        else:
            result = {'error': f'Unknown function: {func_name}'}
        return {
            'role': 'tool',
            'tool_call_id': tc['id'],
            'content': json.dumps(result),
        }

    tasks = [execute_one(idx, tc) for idx, tc in sorted(tool_call_buffers.items())]
    return await asyncio.gather(*tasks)

ツール使用後の最終回答をストリーミングする

ツール呼び出しを実行して結果をメッセージ履歴に追加したら、2回目のストリーミングリクエストを送信してモデルの最終回答を取得します。この応答を直接クライアントへストリーミングします。この2リクエストのパターン(ツール呼び出しを含む初回リクエストと、ツールの結果を含む後続リクエスト)は、エージェントの1ターンを処理する標準的な流れです。どちらのリクエストでも、テキストを UI へストリーミングできます。

async def full_tool_calling_stream(question: str, tools: list):
    messages = [{'role': 'user', 'content': question}]

    # First request: may produce tool calls
    tool_buffers = {}
    text1 = ''
    async for event_type, data in stream_with_tools(messages, tools):
        if event_type == 'text':
            text1 += data
            yield data  # stream partial text if any
        elif event_type == 'tool_calls':
            tool_buffers = data

    if tool_buffers:
        # Execute tools, then get final streaming answer
        tool_results = await execute_tool_calls_parallel(tool_buffers)
        messages += [{  # assistant tool call message
            'role': 'assistant',
            'tool_calls': [
                {'id': tc['id'], 'type': 'function',
                 'function': {'name': tc['name'], 'arguments': tc['arguments']}}
                for tc in tool_buffers.values()
            ]
        }] + tool_results

        # Second request: final answer streams directly
        async for token in token_stream(messages):  # from earlier lesson
            yield token

UI にツール呼び出しの進行状況を表示する

ツール呼び出しの完了を待っている間、ユーザーにはエージェントが何を実行しているかが分かるようにします。ツールを実行する前に、呼び出す関数とその引数を示すステータスイベントをクライアントへストリーミングします。実行後には完了ステータスをストリーミングします。この透明性によって、ユーザーが感じる応答性が大幅に向上し、予期しないツール使用のデバッグにも役立ちます。

import json

async def stream_with_progress(question, tools):
    messages = [{'role': 'user', 'content': question}]
    tool_buffers = {}

    async for event_type, data in stream_with_tools(messages, tools):
        if event_type == 'tool_calls':
            tool_buffers = data

    for tc in tool_buffers.values():
        args = json.loads(tc['arguments'])
        yield f'data: {json.dumps({"type": "tool_start", "function": tc["name"], "args": args})}\n\n'

        result = tools_registry.get(tc['name'], lambda **kw: {})(** args)

        yield f'data: {json.dumps({"type": "tool_done", "function": tc["name"]})}\n\n'

    # Then stream final answer...

ツール呼び出しストリームのエラー処理

ツールの実行は失敗する可能性があります。API がエラーを返したり、関数が例外を発生させたり、JSON の解析に失敗したりすることがあります。ツールの実行では必ず例外を捕捉し、構造化されたエラー応答をモデルに返します。モデルはそれを受けて、異なる引数で再試行するか、代替ツールを呼び出すか、要求された操作が失敗したことをユーザーに説明できます。捕捉されないツール例外によってストリーミングループがクラッシュしないようにしてください。

def safe_execute_tool(func_name: str, args: dict) -> str:
    try:
        if func_name not in tools_registry:
            return json.dumps({'error': f'Function {func_name!r} not found'})
        result = tools_registry[func_name](**args)
        return json.dumps(result)
    except TypeError as e:
        return json.dumps({'error': f'Invalid arguments: {str(e)}'})
    except Exception as e:
        return json.dumps({'error': f'Execution failed: {str(e)}'})

# Tool result message with error handled
tool_message = {
    'role': 'tool',
    'tool_call_id': tc['id'],
    'content': safe_execute_tool(tc['name'], json.loads(tc['arguments'])),
}

エージェントにおけるストリーミングと非ストリーミングの比較

エージェント型アプリケーションでは、ストリーミングによって実装は複雑になりますが、UX に大きな価値がもたらされます。ストリーミングを使わない場合、10〜30秒かかる可能性のある複数ステップのツール呼び出しループの間、ユーザーには何も表示されません。ストリーミングを使えば、中間テキストやツール呼び出しの通知、最終回答がトークン単位で表示されます。コードが複雑になるとしても、インタラクティブなアプリケーションでは通常その価値があります。一方、自律的に実行するバックグラウンドエージェントでは、より単純なコードにするため非ストリーミングを使用できます。

理解度チェック

このレッスンで学んだ、ストリーミングされた応答におけるツール呼び出しについての理解度を確認します。

レッスンのまとめ

このレッスンでは、ストリーミングされた応答ではツール呼び出しの引数がJSON の断片として届くため、解析する前にツール呼び出しのインデックスごとにバッファリングする必要があること、ストリームの完了後にツールを実行し、その後に2回目のストリーミングリクエストを送って最終回答を取得すること、モデルが複数の関数を同時に呼び出す場合は asyncio.gather による並列実行で遅延を最小限に抑えられることを学びました。エージェントループのクラッシュを防ぐため、ツール実行時のエラーは必ず捕捉してください。次は、API コストを削減するための応答キャッシュを実装します。

よくある質問

「ストリーミングレスポンスのツール呼び出しを処理する」レッスンは無料ですか?

はい。「ストリーミングレスポンスのツール呼び出しを処理する」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、AI Engineering Academyコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 AI Engineering Academyコースには全4レッスンが含まれています。

「ストリーミングレスポンスのツール呼び出しを処理する」で何を学びますか?

関数呼び出しの引数がトークン単位で届くストリーミングレスポンスを解析し、JSONの断片をバッファリングして、呼び出しが完了したときだけツールを実行します。 ブラウザで直接実行するハンズオンコードでAI Engineering Academyを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。

AI Engineering Academyを始めるのに経験は必要ですか?

事前経験は必要ありません。CoddyKitのAI Engineering Academyは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン4/4です。

「ストリーミングレスポンスのツール呼び出しを処理する」レッスンにはどのくらい時間がかかりますか?

ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。

このAI Engineering Academyレッスンでコードを書いて実行できますか?

はい。すべてのAI Engineering Academyレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。

このコースのすべてのレッスン

  1. トークンストリーミングを理解する
  2. Python SDKでストリームを利用する
  3. Server-Sent EventsによるFastAPIストリーミング
  4. ストリーミングレスポンスのツール呼び出しを処理する
← AI Engineering Academyに戻る