0Pricing
AI Agents · Lección

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

  1. Python asíncrono para desarrolladores de agentes
  2. Colas de eventos y message brokers
  3. Ejecución paralela no bloqueante de herramientas
  4. Frameworks asíncronos para agentes: LangChain y más
← Volver a AI Agents