Обработка вызовов инструментов в потоковых ответах
Разбирайте потоковые ответы, в которых аргументы вызова функции поступают по одному токену, накапливайте фрагменты JSON и запускайте инструмент только после завершения вызова.
«Обработка вызовов инструментов в потоковых ответах» — бесплатный урок 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 — локальная установка не требуется.
Все уроки этого курса
- Как работает потоковая передача токенов
- Работа с потоками через Python SDK
- Потоковая передача в FastAPI с событиями, отправляемыми сервером
- Обработка вызовов инструментов в потоковых ответах