0Pricing
AI Agents · Lektion

Nicht blockierende parallele Tool-Ausführung

asyncio.gather() zum gleichzeitigen Ausführen mehrerer Tools.

Nicht blockierende parallele Tool-Ausführung ist eine kostenlose AI Agents-Lektion auf CoddyKit. Dies ist Lektion 3 von 4. Du kannst die komplette Lektion unten kostenlos lesen – dann übst du sie direkt im Browser mit einem integrierten Code-Editor und einem KI-Tutor rund um die Uhr. Sie ist Teil des AI Agents-Lernpfads, und dein Fortschritt wird über Web und CoddyKit-App synchronisiert. Der AI Agents-Kurs umfasst insgesamt 4 Lektionen.

Warum parallele Tool-Ausführung?

Wenn ein Agent Ergebnisse von mehreren unabhängigen Tools benötigt, kostet die sequenzielle Ausführung unnötig Zeit. Wenn jedes Tool 500 ms benötigt, dauern drei Tools sequenziell 1,5 s, parallel jedoch nur 500 ms – eine Beschleunigung um das Dreifache.

asyncio.gather() für parallele Aufrufe

asyncio.gather() führt mehrere Coroutinen nebenläufig aus und gibt alle Ergebnisse in der ursprünglichen Reihenfolge zurück. Es ist das wichtigste Werkzeug für die parallele Tool-Ausführung durch Agenten.

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())

Umgang mit individuellen Tool-Fehlern

Mit return_exceptions=True führt der Fehler eines einzelnen Tools nicht zum Abbruch aller parallelen Aufrufe. Jedes Ergebnis ist entweder ein Wert oder eine Exception – prüfen Sie sie einzeln.

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())

Parallele Ergebnisse zusammenführen

Führen Sie die Ergebnisse nach der parallelen Ausführung zu einem einzigen Kontext für das LLM zusammen. Kennzeichnen Sie eindeutig, aus welcher Quelle die einzelnen Informationen stammen.

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

Timeout für einzelne Tools

Verwenden Sie asyncio.wait_for(), um für einzelne Tool-Aufrufe ein Timeout festzulegen. Ein langsames Tool sollte den gesamten parallelen Batch nicht unbegrenzt blockieren.

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())

Dynamische Auswahl paralleler Tools

Ein Agent kann anhand der Abfrage dynamisch entscheiden, welche Tools parallel ausgeführt werden. Erstellen Sie einen Dispatcher, der Tool-Namen asynchronen Funktionen zuordnet und die ausgewählten Funktionen nebenläufig ausführt.

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)

Parallelität mit Semaphoren begrenzen

Zu viele parallele Tool-Aufrufe können API-Ratenlimits überschreiten oder einen Dienst überlasten. Verwenden Sie asyncio.Semaphore, um die Anzahl der gleichzeitig ausgeführten Tool-Aufrufe zu begrenzen.

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')

Streaming von Teilergebnissen

Mit asyncio.as_completed() verarbeiten Sie Tool-Ergebnisse, sobald sie eintreffen, statt auf den Abschluss aller Aufrufe zu warten. Zeigen Sie Benutzern Teilergebnisse sofort an.

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'))

Parallele Tool-Aufrufe in LangChain

LangChain unterstützt parallele Tool-Aufrufe nativ, wenn das LLM mehrere Tool-Aufrufe in einer Antwort zurückgibt. Verarbeiten Sie diese zur Effizienzsteigerung mit 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')

Deduplizierung von Ergebnissen

Bei parallelen Suchen können verschiedene Tools sich überschneidende Ergebnisse zurückgeben. Entfernen Sie Duplikate, bevor Sie die Ergebnisse dem LLM präsentieren, damit dieselben Informationen nicht mehrfach im Kontext erscheinen.

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')

Beschleunigung durch Parallelisierung messen

Messen Sie die tatsächliche Beschleunigung durch die Parallelisierung. Vergleichen Sie die sequenzielle mit der parallelen Ausführungszeit, um den Nutzen zu quantifizieren und die zusätzliche Komplexität zu rechtfertigen.

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))

Wissenscheck: Parallele Ausführung

Testen Sie Ihr Verständnis der nicht blockierenden parallelen Tool-Ausführung.

Zusammenfassung: Parallele Tool-Ausführung

Die parallele Tool-Ausführung mit asyncio reduziert die Latenz von Agenten erheblich. Verwenden Sie asyncio.gather() für nebenläufige Aufrufe, return_exceptions=True für Fehlertoleranz, asyncio.wait_for() für Tool-spezifische Timeouts, Semaphore zur Begrenzung der Aufrufrate und asyncio.as_completed() für das Streaming von Teilergebnissen. Entfernen Sie zusammengeführte Duplikate stets, damit der Kontext für das LLM übersichtlich bleibt.

Häufig gestellte Fragen

Ist die Lektion „Nicht blockierende parallele Tool-Ausführung“ kostenlos?

Ja — der vollständige Text von „Nicht blockierende parallele Tool-Ausführung“ ist hier im Web kostenlos zu lesen. Um sie interaktiv zu üben (integrierter Code-Editor und 24/7 KI-Tutor) und den Rest des AI Agents-Kurses freizuschalten, upgrade auf CoddyKit PRO. Der AI Agents-Kurs umfasst insgesamt 4 Lektionen.

Was lerne ich in „Nicht blockierende parallele Tool-Ausführung“?

asyncio.gather() zum gleichzeitigen Ausführen mehrerer Tools. Du übst AI Agents mit praktischem Code, den du direkt im Browser ausführst, und ein 24/7 KI-Tutor beantwortet deine Fragen während du die Lektion bearbeitest.

Brauche ich Erfahrung, um AI Agents zu starten?

Keine Vorkenntnisse erforderlich. AI Agents auf CoddyKit ist für Anfänger bis fortgeschrittene Lernende strukturiert, sodass du hier starten oder von Anfang an beginnen und in deinem eigenen Tempo voranschreiten kannst. Dies ist Lektion 3 von 4.

Wie lange dauert die Lektion „Nicht blockierende parallele Tool-Ausführung“?

Die meisten CoddyKit-Lektionen dauern etwa 5–10 Minuten. Jede ist kompakt und interaktiv, sodass du stetig Fortschritte machst und genau dort weitermachst, wo du aufgehört hast – im Web und in der App.

Kann ich in dieser AI Agents-Lektion Code schreiben und ausführen?

Ja. Jede AI Agents-Lektion enthält einen integrierten Code-Editor, sodass du echten Code direkt in deinem Browser schreibst und ausführst und sofort KI-Feedback erhältst — ohne lokale Einrichtung erforderlich.

Alle Lektionen in diesem Kurs

  1. Asynchrones Python für Agent-Entwickler
  2. Ereigniswarteschlangen und Message Broker
  3. Nicht blockierende parallele Tool-Ausführung
  4. Asynchrone Agent-Frameworks: LangChain und darüber hinaus
← Zurück zu AI Agents