Tekoälyagentit · Oppitunti

Työkalujen rinnakkainen suoritus ilman estävää odotusta

asyncio.gather() useiden työkalujen samanaikaiseen suorittamiseen.

Oppitunti 3/413 vaihetta

Työkalujen rinnakkainen suoritus ilman estävää odotusta on ilmainen Tekoälyagentit-oppitunti CoddyKitissä. Tämä on oppitunti 3/4. Voit lukea koko oppitunnin alta ilmaiseksi ja harjoitella sen jälkeen käytännössä selaimessa sisäänrakennetulla koodieditorilla ja ympäri vuorokauden käytettävissä olevan tekoälytuutorin avulla. Oppitunti kuuluu Tekoälyagentit-oppimispolkuun, ja edistymisesi synkronoituu verkon ja CoddyKit-sovelluksen välillä. Tekoälyagentit-kurssilla on yhteensä 4 oppituntia.

Miksi suorittaa työkaluja rinnakkain?

Kun agentti tarvitsee tuloksia useista toisistaan riippumattomista työkaluista, niiden suorittaminen peräkkäin tuhlaa aikaa. Jos kunkin työkalun suoritus kestää 500 ms, kolme työkalua kestää peräkkäin suoritettuna 1,5 sekuntia mutta rinnakkain vain 500 ms – nopeutus on kolminkertainen.

asyncio.gather() rinnakkaisiin kutsuihin

asyncio.gather() suorittaa useita korutiineja samanaikaisesti ja palauttaa kaikki tulokset järjestyksessä. Se on ensisijainen työkalu agenttien työkalujen rinnakkaiseen suorittamiseen.

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

Yksittäisten työkalujen virheiden käsittely

Asetuksella return_exceptions=True yksittäisen työkalun virhe ei keskeytä kaikkia rinnakkaisia kutsuja. Jokainen tulos on joko arvo tai poikkeus – tarkistakaa ne yksitellen.

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

Rinnakkaisten tulosten yhdistäminen

Yhdistäkää rinnakkaisen suorituksen jälkeen tulokset yhdeksi LLM:n kontekstiksi. Merkitkää selkeästi, mistä kukin tietojen osa on peräisin.

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

Yksittäisten työkalujen aikakatkaisut

Käyttäkää asyncio.wait_for()-funktiota aikakatkaisun lisäämiseen yksittäisiin työkalukutsuihin. Hidas työkalu ei saa estää koko rinnakkaista erää loputtomasti.

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

Rinnakkaisten työkalujen dynaaminen valinta

Agentti voi päättää kyselyn perusteella dynaamisesti, mitkä työkalut suoritetaan rinnakkain. Rakentakaa välittäjä, joka yhdistää työkalujen nimet asynkronisiin funktioihin ja suorittaa valitut työkalut samanaikaisesti.

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)

Rinnakkaisuuden rajoittaminen semaforeilla

Liian monen työkalukutsun suorittaminen rinnakkain voi ylittää API:n nopeusrajoitukset tai kuormittaa palvelua liikaa. Käyttäkää asyncio.Semaphore-semaforia rajoittamaan samanaikaisesti suoritettavien työkalukutsujen määrää.

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

Osittaisten tulosten suoratoisto

Komennon asyncio.as_completed() avulla voitte käsitellä työkalujen tuloksia niiden saapuessa sen sijaan, että odottaisitte kaikkien valmistumista. Näyttäkää käyttäjille osittaiset tulokset heti.

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

Rinnakkaiset työkalukutsut LangChainissa

LangChain tukee työkalukutsujen rinnakkaisuutta suoraan, kun LLM palauttaa useita työkalukutsuja yhdessä vastauksessa. Käsitelkää ne tehokkaasti komennolla 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')

Tulosten kaksoiskappaleiden poistaminen

Rinnakkaisissa hauissa eri työkalut voivat palauttaa päällekkäisiä tuloksia. Poistakaa kaksoiskappaleet ennen tulosten esittämistä LLM:lle, jotta sama tieto ei esiinny kontekstissa useita kertoja.

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

Rinnakkaisuuden nopeutuksen mittaaminen

Mitatkaa rinnakkaistamisesta saatava todellinen nopeutus. Verratkaa peräkkäisen ja rinnakkaisen suorituksen kestoa, jotta voitte määrittää hyödyn ja perustella lisämonimutkaisuuden.

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

Tietotesti: rinnakkainen suoritus

Testatkaa ymmärrystänne estämättömästä työkalujen rinnakkaisesta suorituksesta.

Työkalujen rinnakkaisen suorituksen yhteenveto

Työkalujen rinnakkainen suoritus asyncio-kirjastolla vähentää agentin viivettä huomattavasti. Käyttäkää asyncio.gather()-funktiota samanaikaisiin kutsuihin, return_exceptions=True-asetusta vikasietoisuuteen, asyncio.wait_for()-funktiota työkalukohtaisiin aikakatkaisuihin, Semaphorea nopeusrajoituksiin ja asyncio.as_completed()-funktiota osittaisten tulosten suoratoistoon. Poistakaa yhdistetyistä tuloksista aina kaksoiskappaleet, jotta LLM:n konteksti säilyy selkeänä.

Aloita maksutta

Opi Tekoälyagentit tekoälytuutorin avulla — ilmaiseksi

Kirjoita ja suorita oikeaa koodia selaimessa, saa välitöntä apua tekoälytuutorilta ympäri vuorokauden ja jatka siitä, mihin jäit, verkossa tai sovelluksessa.

Kurssit
60
Oppitunnit
239

Usein kysytyt kysymykset

Onko oppitunti ”Työkalujen rinnakkainen suoritus ilman estävää odotusta” ilmainen?

Kyllä – oppitunnin ”Työkalujen rinnakkainen suoritus ilman estävää odotusta” koko tekstin voi lukea täällä verkossa ilmaiseksi. Jos haluat harjoitella interaktiivisesti sisäänrakennetulla koodieditorilla ja ympäri vuorokauden käytettävissä olevan tekoälytuutorin avulla sekä avata koko Tekoälyagentit-kurssin, päivitä CoddyKit PROhon. Tekoälyagentit-kurssilla on yhteensä 4 oppituntia.

Mitä opin oppitunnilla ”Työkalujen rinnakkainen suoritus ilman estävää odotusta”?

asyncio.gather() useiden työkalujen samanaikaiseen suorittamiseen. Harjoittelet Tekoälyagentit-aihetta koodilla, jonka suoritat suoraan selaimessa. Ympäri vuorokauden käytettävissä oleva tekoälytuutori vastaa kysymyksiisi oppitunnin aikana.

Tarvitsenko kokemusta aloittaakseni Tekoälyagentit-opiskelun?

Aiempi kokemus ei ole tarpeen. CoddyKitin Tekoälyagentit-oppimispolku sopii vasta-alkajista edistyneisiin, joten voit aloittaa tästä tai alusta ja edetä omaan tahtiisi. Tämä on oppitunti 3/4.

Kuinka kauan ”Työkalujen rinnakkainen suoritus ilman estävää odotusta”-oppitunnin suorittaminen kestää?

Useimmat CoddyKitin oppitunnit kestävät noin 5–10 minuuttia. Jokainen oppitunti on lyhyt ja interaktiivinen, joten edistyt tasaisesti ja voit jatkaa siitä, mihin jäit – sekä verkossa että sovelluksessa.

Voinko kirjoittaa ja suorittaa koodia tällä Tekoälyagentit-oppitunnilla?

Kyllä. Jokainen Tekoälyagentit-oppitunti sisältää sisäänrakennetun koodieditorin, joten voit kirjoittaa ja suorittaa oikeaa koodia suoraan selaimessa ja saada välitöntä palautetta tekoälyltä – paikallista asennusta ei tarvita.

Kaikki tämän kurssin oppitunnit

  1. Asynkroninen Python agenttikehittäjille
  2. Tapahtumajonot ja viestinvälittäjät
  3. Työkalujen rinnakkainen suoritus ilman estävää odotusta
  4. Asynkroniset agenttikehykset: LangChain ja muut
← Takaisin: Tekoälyagentit