AI Engineering Academy · Урок

Обработка вызовов инструментов в потоковых ответах

Разбирайте потоковые ответы, в которых аргументы вызова функции поступают по одному токену, накапливайте фрагменты JSON и запускайте инструмент только после завершения вызова.

Урок 4 из 413 шагов

«Обработка вызовов инструментов в потоковых ответах» — бесплатный урок AI Engineering Academy на CoddyKit. Это урок 4 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения AI Engineering Academy, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс AI Engineering Academy содержит 4 уроков всего.

Вызовы инструментов в потоках поступают иначе

Когда LLM решает вызвать функцию, структура ответа меняется. Вместо строки content изменение содержит массив tool_calls. Однако в потоковом ответе аргументы вызова функции поступают токен за токеном как частичная строка JSON — полный объект 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

Буферизация аргументов вызовов инструментов

Используйте словарь с индексом вызова инструмента в качестве ключа, чтобы накапливать фрагменты аргументов из каждого фрагмента. Индекс указывает, к какому вызову инструмента относится текущая потоковая передача: модель может вызвать несколько функций в одном ответе. Для каждого фрагмента с ненулевым дельта-значением tool_calls добавляйте фрагмент аргумента в соответствующую запись буфера, определяемую индексом вызова инструмента.

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

Полный цикл потоковой передачи вызовов инструментов

Полный процесс потоковой передачи вызовов инструментов требует цикла многошагового диалога. Первый запрос может вернуть вызовы инструментов; Вы выполняете их и добавляете результаты в историю сообщений, а второй запрос возвращает итоговый текстовый ответ. Этот цикл может повторяться несколько раз, если модель решит выполнить дополнительные вызовы инструментов на основе результатов предыдущих.

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)

Потоковая передача итогового ответа после использования инструментов

После выполнения вызовов инструментов и добавления результатов в историю сообщений отправьте второй потоковый запрос, чтобы получить итоговый ответ модели. Передавайте этот ответ непосредственно клиенту. Такая схема из двух запросов (первоначальный запрос с вызовами инструментов и последующий запрос с результатами их работы) является стандартным циклом хода агента; оба запроса могут передавать текст в интерфейс в потоковом режиме.

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

Отображение хода выполнения вызовов инструментов в интерфейсе

Пользователи должны видеть что делает агент, пока выполняются вызовы инструментов. Перед выполнением инструментов передайте клиенту событие состояния, указывающее, какая функция вызывается и с какими аргументами. После выполнения передайте событие о завершении. Такая прозрачность значительно улучшает воспринимаемую отзывчивость и помогает пользователям находить причины неожиданного использования инструментов.

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'])),
}

Сравнение потокового и непотокового режима для агентов

Для приложений с агентами потоковая передача усложняет реализацию, но значительно улучшает взаимодействие с пользователем. Без потоковой передачи пользователь ничего не видит во время многошагового цикла вызовов инструментов, который может занять 10–30 секунд. При потоковой передаче он видит промежуточный текст, уведомления о вызовах инструментов и итоговый ответ, появляющийся токен за токеном. Дополнительная сложность кода обычно оправдана для интерактивных приложений, однако фоновые агенты, работающие автономно, могут использовать непотоковый режим для более простой реализации.

Быстрая проверка

Проверьте, насколько хорошо Вы поняли работу вызовов инструментов в потоковых ответах из этого урока.

Итоги урока

В этом уроке Вы узнали, что аргументы вызовов инструментов поступают в потоковых ответах как фрагменты JSON и перед разбором должны накапливаться с учётом индекса вызова инструмента; после завершения потока нужно выполнить инструменты, а затем отправить второй потоковый запрос для получения итогового ответа; параллельное выполнение с помощью asyncio.gather уменьшает задержку, когда модель одновременно вызывает несколько функций. Всегда перехватывайте ошибки выполнения инструментов, чтобы предотвратить аварийное завершение цикла агента. Далее мы реализуем кэширование ответов, чтобы сократить расходы на API.

Можно начать бесплатно

Изучай Python с ИИ-репетитором — бесплатно

Пиши и запускай код прямо в браузере, получай мгновенную помощь от ИИ-репетитора 24/7 и продолжи учиться на сайте или в приложении.

Курсы
30
Уроки
120

Часто задаваемые вопросы

Урок «Обработка вызовов инструментов в потоковых ответах» бесплатный?

Да — полный текст урока «Обработка вызовов инструментов в потоковых ответах» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс AI Engineering Academy, подпишись на CoddyKit PRO. Курс AI Engineering Academy содержит 4 уроков всего.

Чему я научусь в уроке «Обработка вызовов инструментов в потоковых ответах»?

Разбирайте потоковые ответы, в которых аргументы вызова функции поступают по одному токену, накапливайте фрагменты JSON и запускайте инструмент только после завершения вызова. Ты практикуешь AI Engineering Academy с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.

Нужен ли мне опыт, чтобы начать AI Engineering Academy?

Предыдущий опыт не требуется. AI Engineering Academy на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 4 из 4.

Сколько времени занимает урок «Обработка вызовов инструментов в потоковых ответах»?

Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.

Можно ли писать и запускать код в этом уроке AI Engineering Academy?

Да. Каждый урок AI Engineering Academy включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.

Все уроки этого курса

  1. Как работает потоковая передача токенов
  2. Работа с потоками через Python SDK
  3. Потоковая передача в FastAPI с событиями, отправляемыми сервером
  4. Обработка вызовов инструментов в потоковых ответах
← Назад к AI Engineering Academy