Ejecución paralela no bloqueante de herramientas
Uso de asyncio.gather() para ejecutar varias herramientas simultáneamente.
Ejecución paralela no bloqueante de herramientas es una lección gratuita de AI Agents en CoddyKit. Esta es la lección 3 de 4. Puedes leer la lección completa abajo gratuitamente — luego la practicas en el navegador con un editor de código integrado y un tutor de IA 24/7. Forma parte de la ruta de aprendizaje de AI Agents, y tu progreso se sincroniza en la web y la app de CoddyKit. El curso de AI Agents incluye 4 lecciones en total.
¿Por qué ejecutar herramientas en paralelo?
Cuando un agente necesita resultados de varias herramientas independientes, ejecutarlas secuencialmente supone una pérdida de tiempo. Si cada herramienta tarda 500 ms, tres herramientas tardan 1,5 s de forma secuencial, pero solo 500 ms en paralelo: una aceleración de 3 veces.
asyncio.gather() para llamadas en paralelo
asyncio.gather() ejecuta varias corrutinas de forma concurrente y devuelve todos los resultados en orden. Es la herramienta principal para ejecutar herramientas de agentes en paralelo.
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())Gestión de fallos individuales de las herramientas
Con return_exceptions=True, el fallo de una sola herramienta no interrumpe todas las llamadas en paralelo. Cada resultado es un valor o una excepción; compruébelos individualmente.
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())Combinar resultados en paralelo
Después de la ejecución en paralelo, combine los resultados en un único contexto para el LLM. Indique claramente el origen de cada fragmento de información.
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'])Tiempo de espera para herramientas individuales
Use asyncio.wait_for() para añadir un tiempo de espera a las llamadas de herramientas individuales. Una herramienta lenta no debe bloquear indefinidamente todo el lote paralelo.
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())Selección dinámica de herramientas en paralelo
Un agente puede decidir dinámicamente qué herramientas ejecutar en paralelo según la consulta. Cree un distribuidor que asigne nombres de herramientas a funciones asíncronas y ejecute simultáneamente las herramientas seleccionadas.
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)Limitar el paralelismo con semáforos
Ejecutar demasiadas llamadas de herramientas en paralelo puede alcanzar los límites de tasa de una API o sobrecargar un servicio. Use asyncio.Semaphore para limitar cuántas llamadas de herramientas se ejecutan simultáneamente.
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')Transmitir resultados parciales
Con asyncio.as_completed(), procese los resultados de las herramientas a medida que llegan en lugar de esperar a que terminen todas. Muestre a los usuarios los resultados parciales de inmediato.
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'))Llamadas de herramientas en paralelo en LangChain
LangChain admite de forma nativa las llamadas de herramientas en paralelo cuando el LLM devuelve varias llamadas de herramientas en una sola respuesta. Gestiónelas con asyncio.gather() para mejorar la eficiencia.
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')Eliminación de duplicados en los resultados
Al ejecutar búsquedas en paralelo, distintas herramientas pueden devolver resultados coincidentes. Elimine los duplicados antes de presentarlos al LLM para evitar que la misma información aparezca varias veces en el contexto.
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')Medir la aceleración del paralelismo
Mida la aceleración real obtenida mediante el paralelismo. Compare el tiempo de ejecución secuencial con el paralelo para cuantificar el beneficio y justificar la complejidad añadida.
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))Comprobación de conocimientos: ejecución en paralelo
Ponga a prueba sus conocimientos sobre la ejecución en paralelo de herramientas sin bloqueo.
Resumen de la ejecución de herramientas en paralelo
La ejecución de herramientas en paralelo con asyncio reduce notablemente la latencia de los agentes. Use asyncio.gather() para llamadas concurrentes, return_exceptions=True para tolerancia a fallos, asyncio.wait_for() para tiempos de espera por herramienta, Semaphore para limitar la tasa y asyncio.as_completed() para transmitir resultados parciales. Elimine siempre los duplicados de los resultados combinados para mantener limpio el contexto del LLM.
Preguntas frecuentes
¿La lección «Ejecución paralela no bloqueante de herramientas» es gratis?
Sí — el texto completo de «Ejecución paralela no bloqueante de herramientas» es gratis para leer aquí en la web. Para practicarla de forma interactiva (editor de código integrado y tutor de IA 24/7) y desbloquear el resto del curso de AI Agents, actualiza a CoddyKit PRO. El curso de AI Agents incluye 4 lecciones en total.
¿Qué aprenderé en «Ejecución paralela no bloqueante de herramientas»?
Uso de asyncio.gather() para ejecutar varias herramientas simultáneamente. Practicas AI Agents con código real que ejecutas directamente en el navegador, y un tutor de IA 24/7 responde tus preguntas mientras trabajas en la lección.
¿Necesito experiencia previa para empezar AI Agents?
No se requiere experiencia previa. AI Agents en CoddyKit está estructurado para principiantes hasta estudiantes avanzados, así que puedes empezar aquí o desde el inicio y avanzar a tu ritmo. Esta es la lección 3 de 4.
¿Cuánto tiempo toma la lección «Ejecución paralela no bloqueante de herramientas»?
La mayoría de las lecciones de CoddyKit toman alrededor de 5–10 minutos. Cada una es compacta e interactiva, así que avanzas constantemente y retomas exactamente por donde dejaste en la web y la app.
¿Puedo escribir y ejecutar código en esta lección de AI Agents?
Sí. Cada lección de AI Agents incluye un editor de código integrado, así que escribes y ejecutas código real directamente en tu navegador y obtienes retroalimentación instantánea de IA — sin configuración local necesaria.
Todas las lecciones de este curso
- Python asíncrono para desarrolladores de agentes
- Colas de eventos y message brokers
- Ejecución paralela no bloqueante de herramientas
- Frameworks asíncronos para agentes: LangChain y más