Неблокирующее параллельное выполнение инструментов
Использование asyncio.gather() для одновременного запуска нескольких инструментов.
«Неблокирующее параллельное выполнение инструментов» — бесплатный урок AI Agents на CoddyKit. Это урок 3 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения AI Agents, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс AI Agents содержит 4 уроков всего.
Зачем выполнять инструменты параллельно?
Когда агенту нужны результаты нескольких независимых инструментов, последовательное выполнение приводит к потере времени. Если каждый инструмент работает 500 мс, три инструмента выполняются последовательно за 1,5 с, а параллельно — всего за 500 мс, что даёт ускорение в 3 раза.
asyncio.gather() для параллельных вызовов
asyncio.gather() выполняет несколько сопрограмм одновременно и возвращает все результаты в исходном порядке. Это основной инструмент для параллельного выполнения вызовов инструментов агентом.
import asyncio
import time
async def search_web(query: str) -> list:
await asyncio.sleep(0.5) # Simulate 500ms web search
return [f'Web result for: {query}']
async def search_database(query: str) -> list:
await asyncio.sleep(0.3) # Simulate 300ms DB query
return [f'DB result for: {query}']
async def get_weather(location: str) -> dict:
await asyncio.sleep(0.4) # Simulate 400ms API call
return {'location': location, 'temp': '22C'}
async def run_parallel():
start = time.perf_counter()
# Sequential: 0.5 + 0.3 + 0.4 = 1.2s
# Parallel: max(0.5, 0.3, 0.4) = 0.5s
web_results, db_results, weather = await asyncio.gather(
search_web('Python async'),
search_database('Python async'),
get_weather('New York')
)
elapsed = (time.perf_counter() - start) * 1000
print(f'Completed in {elapsed:.0f}ms (parallel)')
return web_results, db_results, weather
asyncio.run(run_parallel())Обработка отдельных сбоев инструментов
При использовании return_exceptions=True сбой одного инструмента не прерывает все параллельные вызовы. Каждый результат представляет собой либо значение, либо исключение — проверяйте их по отдельности.
import asyncio
async def tool_that_fails(name: str):
await asyncio.sleep(0.2)
if name == 'flaky_api':
raise ConnectionError(f'{name}: service unavailable')
return f'{name}: success'
async def parallel_with_fault_tolerance():
tool_names = ['web_search', 'flaky_api', 'database', 'weather_api']
coros = [tool_that_fails(name) for name in tool_names]
results = await asyncio.gather(*coros, return_exceptions=True)
tool_results = {}
errors = {}
for name, result in zip(tool_names, results):
if isinstance(result, Exception):
errors[name] = str(result)
print(f'Tool {name} FAILED: {result}')
else:
tool_results[name] = result
print(f'Tool {name} OK: {result}')
print(f'\nSucceeded: {len(tool_results)}/{len(tool_names)}')
print('Errors:', errors)
return tool_results, errors
asyncio.run(parallel_with_fault_tolerance())Объединение параллельных результатов
После параллельного выполнения объедините результаты в единый контекст для LLM. Чётко указывайте источник каждого фрагмента информации.
import asyncio
import json
async def parallel_research(query: str) -> dict:
web_task = search_web(query)
db_task = search_database(query)
weather_task = get_weather('New York')
results = await asyncio.gather(
web_task, db_task, weather_task,
return_exceptions=True
)
context_parts = []
sources_used = []
tool_names = ['web_search', 'database', 'weather']
for name, result in zip(tool_names, results):
if isinstance(result, Exception):
context_parts.append(f'[{name}]: unavailable ({result})')
else:
context_parts.append(f'[{name}]: {json.dumps(result)}')
sources_used.append(name)
combined_context = '\n'.join(context_parts)
return {
'query': query,
'context': combined_context,
'sources': sources_used
}
result = asyncio.run(parallel_research('Python performance tips'))
print('Sources used:', result['sources'])Время ожидания для отдельных инструментов
Используйте asyncio.wait_for(), чтобы добавить время ожидания к отдельным вызовам инструментов. Медленный инструмент не должен бесконечно блокировать всю параллельную пакетную операцию.
import asyncio
async def slow_tool(name: str) -> str:
await asyncio.sleep(10) # Very slow
return f'{name} result'
async def tool_with_timeout(coro, tool_name: str, timeout_seconds: float):
try:
result = await asyncio.wait_for(coro, timeout=timeout_seconds)
return result
except asyncio.TimeoutError:
return f'TIMEOUT: {tool_name} exceeded {timeout_seconds}s'
except Exception as e:
return f'ERROR: {tool_name}: {str(e)}'
async def parallel_with_timeouts():
results = await asyncio.gather(
tool_with_timeout(search_web('query'), 'web_search', 2.0),
tool_with_timeout(slow_tool('slow_api'), 'slow_api', 1.0),
tool_with_timeout(search_database('query'), 'database', 2.0)
)
for result in results:
print(result)
asyncio.run(parallel_with_timeouts())Динамический выбор инструментов для параллельного выполнения
Агент может динамически выбирать инструменты для параллельного запуска в зависимости от запроса. Создайте диспетчер, который сопоставляет имена инструментов с асинхронными функциями и одновременно запускает выбранные функции.
import asyncio
TOOL_REGISTRY = {
'web_search': search_web,
'database': search_database,
'weather': get_weather
}
async def execute_tools_parallel(tool_calls: list) -> dict:
'''
tool_calls: list of {'name': str, 'args': dict}
'''
tasks = {}
for call in tool_calls:
tool_name = call['name']
args = call.get('args', {})
fn = TOOL_REGISTRY.get(tool_name)
if fn:
# Get first positional arg (simplified)
first_arg = next(iter(args.values()), '') if args else ''
tasks[tool_name] = fn(first_arg)
else:
print(f'Unknown tool: {tool_name}')
if not tasks:
return {}
results = await asyncio.gather(*tasks.values(), return_exceptions=True)
return {
name: result
for name, result in zip(tasks.keys(), results)
}
tool_calls = [
{'name': 'web_search', 'args': {'query': 'async Python'}},
{'name': 'weather', 'args': {'location': 'London'}}
]
results = asyncio.run(execute_tools_parallel(tool_calls))
print('Results:', results)Ограничение параллелизма с помощью семафоров
Слишком большое количество параллельных вызовов инструментов может привести к превышению ограничений частоты запросов API или перегрузке сервиса. Используйте asyncio.Semaphore, чтобы ограничить число одновременно выполняющихся вызовов инструментов.
import asyncio
MAX_CONCURRENT = 3
semaphore = asyncio.Semaphore(MAX_CONCURRENT)
async def rate_limited_tool(tool_fn, *args):
async with semaphore: # Blocks if MAX_CONCURRENT calls already running
return await tool_fn(*args)
async def process_many_queries(queries: list) -> list:
print(f'Processing {len(queries)} queries with max {MAX_CONCURRENT} concurrent')
tasks = [rate_limited_tool(search_web, q) for q in queries]
results = await asyncio.gather(*tasks, return_exceptions=True)
successes = [r for r in results if not isinstance(r, Exception)]
print(f'Completed: {len(successes)}/{len(queries)}')
return results
queries = [f'query-{i}' for i in range(10)]
results = asyncio.run(process_many_queries(queries))
print('Done')Потоковая передача частичных результатов
С помощью asyncio.as_completed() обрабатывайте результаты инструментов по мере их поступления, а не ждите завершения всех вызовов. Показывайте пользователям частичные результаты сразу.
import asyncio
async def search_slow(query: str) -> dict:
await asyncio.sleep(1.0)
return {'source': 'slow_db', 'results': [f'Slow result for: {query}']}
async def search_fast(query: str) -> dict:
await asyncio.sleep(0.2)
return {'source': 'fast_cache', 'results': [f'Fast result for: {query}']}
async def process_as_available(query: str):
coros = [
search_fast(query),
search_slow(query),
search_web(query),
search_database(query)
]
tasks = [asyncio.create_task(c) for c in coros]
partial_results = []
print('Processing results as they arrive:')
for future in asyncio.as_completed(tasks):
result = await future
partial_results.append(result)
print(f' Got result {len(partial_results)}: {result}')
# In a real agent: stream this to the user interface
return partial_results
asyncio.run(process_as_available('machine learning'))Параллельные вызовы инструментов в LangChain
LangChain изначально поддерживает параллельные вызовы инструментов, когда LLM возвращает несколько вызовов инструментов в одном ответе. Для эффективности обрабатывайте их с помощью asyncio.gather().
import asyncio
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage
llm = ChatOpenAI(model='gpt-4o-mini', api_key='sk-...')
TOOL_EXECUTORS = {
'search_web': lambda args: search_web(args.get('query', '')),
'search_database': lambda args: search_database(args.get('query', '')),
'get_weather': lambda args: get_weather(args.get('location', 'New York'))
}
async def handle_parallel_tool_calls(response_message) -> list:
if not response_message.tool_calls:
return []
tasks = []
tool_call_ids = []
for tool_call in response_message.tool_calls:
executor = TOOL_EXECUTORS.get(tool_call['name'])
if executor:
tasks.append(executor(tool_call['args']))
tool_call_ids.append(tool_call['id'])
results = await asyncio.gather(*tasks, return_exceptions=True)
return [
{'tool_call_id': tc_id, 'result': r}
for tc_id, r in zip(tool_call_ids, results)
]
print('Parallel LangChain tool execution defined')Удаление дубликатов результатов
При выполнении параллельного поиска разные инструменты могут возвращать пересекающиеся результаты. Удаляйте дубликаты до передачи результатов LLM, чтобы одна и та же информация не появлялась в контексте несколько раз.
import hashlib
def deduplicate_results(all_results: list) -> list:
seen_hashes = set()
unique_results = []
for result in all_results:
content = str(result)
content_hash = hashlib.md5(content.encode()).hexdigest()
if content_hash not in seen_hashes:
seen_hashes.add(content_hash)
unique_results.append(result)
return unique_results
def merge_parallel_results(tool_results: dict) -> list:
all_items = []
for tool_name, results in tool_results.items():
if isinstance(results, Exception):
continue
if isinstance(results, list):
for item in results:
if isinstance(item, dict):
item['source'] = tool_name
all_items.append(item)
else:
all_items.append({'source': tool_name, 'data': results})
return deduplicate_results(all_items)
sample_results = {
'web': ['Result A', 'Result B'],
'db': ['Result B', 'Result C'] # Result B is duplicate
}
merged = merge_parallel_results(sample_results)
print(f'Before: {sum(len(v) for v in sample_results.values())} items')
print(f'After dedup: {len(merged)} items')Измерение ускорения при параллельном выполнении
Измеряйте фактическое ускорение от распараллеливания. Сравнивайте время последовательного и параллельного выполнения, чтобы оценить пользу и обосновать дополнительную сложность.
import asyncio
import time
async def measure_speedup(tools_and_args: list):
# Sequential timing
seq_start = time.perf_counter()
seq_results = []
for fn, args in tools_and_args:
result = await fn(*args)
seq_results.append(result)
seq_time = (time.perf_counter() - seq_start) * 1000
# Parallel timing
par_start = time.perf_counter()
par_results = await asyncio.gather(*[fn(*args) for fn, args in tools_and_args])
par_time = (time.perf_counter() - par_start) * 1000
speedup = seq_time / par_time if par_time > 0 else 0
print(f'Sequential: {seq_time:.0f}ms')
print(f'Parallel: {par_time:.0f}ms')
print(f'Speedup: {speedup:.1f}x')
return speedup
tools = [
(search_web, ('query',)),
(search_database, ('query',)),
(get_weather, ('London',))
]
asyncio.run(measure_speedup(tools))Проверка знаний: параллельное выполнение
Проверьте своё понимание неблокирующего параллельного выполнения инструментов.
Итоги по параллельному выполнению инструментов
Параллельное выполнение инструментов с помощью asyncio значительно сокращает задержку агента. Используйте asyncio.gather() для одновременных вызовов, return_exceptions=True для устойчивости к сбоям, asyncio.wait_for() для ограничения времени ожидания отдельных инструментов, Semaphore для ограничения частоты запросов и asyncio.as_completed() для потоковой передачи частичных результатов. Всегда удаляйте дубликаты объединённых результатов, чтобы контекст LLM оставался чистым.
Часто задаваемые вопросы
Урок «Неблокирующее параллельное выполнение инструментов» бесплатный?
Да — полный текст урока «Неблокирующее параллельное выполнение инструментов» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс AI Agents, подпишись на CoddyKit PRO. Курс AI Agents содержит 4 уроков всего.
Чему я научусь в уроке «Неблокирующее параллельное выполнение инструментов»?
Использование asyncio.gather() для одновременного запуска нескольких инструментов. Ты практикуешь AI Agents с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.
Нужен ли мне опыт, чтобы начать AI Agents?
Предыдущий опыт не требуется. AI Agents на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 3 из 4.
Сколько времени занимает урок «Неблокирующее параллельное выполнение инструментов»?
Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.
Можно ли писать и запускать код в этом уроке AI Agents?
Да. Каждый урок AI Agents включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.
Все уроки этого курса
- Асинхронный Python для разработчиков агентов
- Очереди событий и брокеры сообщений
- Неблокирующее параллельное выполнение инструментов
- Асинхронные фреймворки агентов: LangChain и не только