面向智能体开发者的异步 Python
asyncio 基础、async def、await、事件循环——理解异步模型
面向智能体开发者的异步 Python 是 CoddyKit 上的免费 AI Agents 课时。 这是第 1 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 AI Agents 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 AI Agents 课程共包含 4 节课。
为什么智能体需要异步处理
智能体会发起许多受 I/O 限制的调用,例如 LLM 接口、网络请求和数据库查询。同步代码在这些调用完成期间只能空闲等待。异步代码可以在等待期间执行其他工作,从而大幅提高吞吐量。
异步基础:async def 与 await
async def 用于声明协程函数。await 会暂停执行,直到等待的操作完成,同时让事件循环在此期间运行其他协程。
import asyncio
async def fetch_data(source: str) -> str:
print(f'Starting fetch from {source}')
await asyncio.sleep(1) # Simulates a network call
print(f'Finished fetch from {source}')
return f'Data from {source}'
async def main():
# Sequential: takes 2 seconds total
result1 = await fetch_data('source-A')
result2 = await fetch_data('source-B')
print('Sequential results:', result1, result2)
asyncio.run(main())
# asyncio.run() starts the event loop and runs main()
# It is the entry point for async programs事件循环
事件循环是 asyncio 的核心。它管理协程和 I/O 回调队列,并在它们准备就绪时运行它们。所有异步代码都在单个线程上的事件循环内运行。
import asyncio
async def task_a():
print('Task A: start')
await asyncio.sleep(2)
print('Task A: done')
async def task_b():
print('Task B: start')
await asyncio.sleep(1)
print('Task B: done')
async def main():
# asyncio.gather runs both tasks concurrently
# Total time: ~2 seconds (not 3)
await asyncio.gather(task_a(), task_b())
print('Both tasks complete')
# Expected output order:
# Task A: start
# Task B: start
# Task B: done <- after 1s
# Task A: done <- after 2s
# Both tasks complete
asyncio.run(main())协程与线程
协程具有协作性:它们通过 await 显式让出控制权。线程具有抢占性:OS 可以随时在线程之间切换。协程更轻量,并且在处理 I/O 时不存在 GIL 问题,也更容易理解和推理。
import asyncio
import threading
import time
# Thread approach: multiple OS threads
def thread_worker(name):
print(f'Thread {name}: start')
time.sleep(1) # Blocks the thread
print(f'Thread {name}: done')
threads = [threading.Thread(target=thread_worker, args=(i,)) for i in range(3)]
for t in threads:
t.start()
for t in threads:
t.join()
print('---')
# Coroutine approach: single-threaded event loop
async def coro_worker(name):
print(f'Coro {name}: start')
await asyncio.sleep(1) # Suspends, does NOT block other coroutines
print(f'Coro {name}: done')
async def main():
await asyncio.gather(*[coro_worker(i) for i in range(3)])
asyncio.run(main())
# Both approaches run 3 tasks in ~1 second total, but coroutines use one threadasyncio.run() 入口点
asyncio.run() 会创建新的事件循环,运行给定的协程直到完成,然后关闭该循环。这是 Python 3.7 及更高版本中异步程序的标准入口点。
import asyncio
async def agent_main():
print('Agent starting')
# All agent async work goes here
results = await asyncio.gather(
asyncio.sleep(0.1), # Simulated LLM call
asyncio.sleep(0.1), # Simulated DB query
)
print('Agent done')
return 'complete'
# Run the agent
result = asyncio.run(agent_main())
print('Result:', result)
# WRONG: calling asyncio.run() inside an already-running event loop
# In Jupyter notebooks, use: await agent_main() directly
# Or: nest_asyncio.apply() then asyncio.run()常见错误:在异步上下文中执行阻塞操作
不要在异步代码中调用阻塞函数(time.sleep、requests.get、同步文件 I/O)。这会阻塞整个事件循环,使所有并发执行都无法进行。
import asyncio
import time
import httpx
# WRONG: blocks the event loop
async def bad_agent_step():
time.sleep(2) # Blocks all other coroutines
# requests.get(url) # Also blocks - do NOT use requests in async code
return 'done'
# RIGHT: use async equivalents
async def good_agent_step():
await asyncio.sleep(2) # Suspends, other coroutines can run
async with httpx.AsyncClient() as client:
response = await client.get('https://api.example.com/data')
return response.text
# For CPU-intensive work: use run_in_executor
import concurrent.futures
async def cpu_intensive_step(data: str):
loop = asyncio.get_event_loop()
with concurrent.futures.ProcessPoolExecutor() as pool:
result = await loop.run_in_executor(pool, expensive_cpu_fn, data)
return result
def expensive_cpu_fn(data):
# CPU-bound work runs in separate process
return data.upper()
print('Blocking vs non-blocking patterns demonstrated')忘记 await:隐蔽的错误
忘记 await 不会引发错误,而是返回协程对象,而不是操作结果。这是一种隐蔽的错误,可能导致后续错误或空结果。
import asyncio
async def get_answer() -> str:
await asyncio.sleep(0.1)
return 'The answer is 42'
async def bad_call():
result = get_answer() # WRONG: forgot await
print(type(result)) # <class 'coroutine'> - not a string!
# Using result as a string here causes AttributeError or wrong behavior
return result
async def good_call():
result = await get_answer() # CORRECT
print(type(result)) # <class 'str'>
return result
async def main():
bad = await bad_call()
print('Bad result:', bad) # coroutine object, not the string
good = await good_call()
print('Good result:', good) # 'The answer is 42'
# Clean up the uncollected coroutine
if asyncio.iscoroutine(bad):
bad.close()
asyncio.run(main())嵌套事件循环问题
在已经运行的事件循环中调用 asyncio.run()(例如在 Jupyter 或 FastAPI 中)会引发 RuntimeError。解决方案是直接使用 await,或者在笔记本中使用 nest_asyncio。
import asyncio
async def my_agent_coroutine():
await asyncio.sleep(0.1)
return 'done'
# In FastAPI or other async frameworks, the event loop is already running
# Use await directly in async endpoints:
async def fastapi_endpoint():
# WRONG inside async context:
# result = asyncio.run(my_agent_coroutine()) # RuntimeError!
# CORRECT: just await
result = await my_agent_coroutine()
return result
# In Jupyter notebooks: install nest_asyncio
# import nest_asyncio
# nest_asyncio.apply()
# Then asyncio.run() works
# Detect if running in event loop:
def run_agent(coro):
try:
loop = asyncio.get_running_loop()
# Already in async context
import concurrent.futures
with concurrent.futures.ThreadPoolExecutor() as pool:
future = pool.submit(asyncio.run, coro)
return future.result()
except RuntimeError:
# No running loop
return asyncio.run(coro)
print('Nested event loop solution defined')asyncio.create_task
使用 asyncio.create_task() 调度协程运行,而无需立即等待它完成。这样您就可以启动多个任务,稍后再等待它们完成。
import asyncio
async def background_job(job_id: int) -> str:
await asyncio.sleep(0.5)
return f'Job {job_id} completed'
async def main():
# Start all tasks without waiting
task1 = asyncio.create_task(background_job(1))
task2 = asyncio.create_task(background_job(2))
task3 = asyncio.create_task(background_job(3))
# Do other work while tasks run
print('Tasks started, doing other work...')
await asyncio.sleep(0.1)
print('Other work done')
# Now wait for all tasks
results = await asyncio.gather(task1, task2, task3)
print('All results:', results)
# Or wait for the first to complete
task_a = asyncio.create_task(background_job(4))
task_b = asyncio.create_task(background_job(5))
done, pending = await asyncio.wait([task_a, task_b], return_when=asyncio.FIRST_COMPLETED)
for t in pending:
t.cancel() # Cancel remaining tasks
print('First result:', done.pop().result())
asyncio.run(main())异步上下文管理器
许多异步库会通过 async with 使用异步上下文管理器。这可以确保异步代码中的连接和资源得到正确的设置与清理。
import asyncio
import httpx
async def fetch_multiple_urls(urls: list) -> list:
# async with ensures the client is properly closed
async with httpx.AsyncClient(timeout=10.0) as client:
# Fetch all URLs concurrently
tasks = [client.get(url) for url in urls]
responses = await asyncio.gather(*tasks, return_exceptions=True)
results = []
for url, response in zip(urls, responses):
if isinstance(response, Exception):
results.append({'url': url, 'error': str(response)})
else:
results.append({'url': url, 'status': response.status_code})
return results
# Async generators for streaming
async def stream_agent_events():
events = ['thinking', 'searching', 'generating', 'done']
for event in events:
await asyncio.sleep(0.2) # Simulate event arrival
yield event
async def consume_stream():
async for event in stream_agent_events():
print(f'Event: {event}')
asyncio.run(consume_stream())异步代码中的错误处理
使用 asyncio.gather(..., return_exceptions=True) 捕获各个任务的失败,而不会中止整个批次。请检查每个结果是否属于异常类型。
import asyncio
async def might_fail(task_id: int) -> str:
await asyncio.sleep(0.1)
if task_id == 2:
raise ValueError(f'Task {task_id} failed')
return f'Task {task_id} succeeded'
async def robust_gather():
tasks = [might_fail(i) for i in range(1, 5)]
# return_exceptions=True: exceptions are returned as values, not raised
results = await asyncio.gather(*tasks, return_exceptions=True)
successes = []
failures = []
for i, result in enumerate(results):
if isinstance(result, Exception):
failures.append({'task': i + 1, 'error': str(result)})
else:
successes.append(result)
print(f'Succeeded: {len(successes)}, Failed: {len(failures)}')
print('Failures:', failures)
return successes, failures
asyncio.run(robust_gather())知识检查:异步 Python
请测试您对用于智能体开发的异步 Python 的理解。
异步 Python 总结
面向智能体开发者的异步 Python 关键规则:对所有受 I/O 限制的操作使用 async def/await;不要在异步上下文中调用阻塞函数;使用 asyncio.gather() 执行并行操作;使用 return_exceptions=True 实现容错的并行调用;使用 async with 管理资源;在进程池执行器中运行 CPU 密集型工作。
常见问题解答
「面向智能体开发者的异步 Python」课时是免费的吗?
是的 — 「面向智能体开发者的异步 Python」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 AI Agents 课程的其余内容,请升级到 CoddyKit PRO。 AI Agents 课程共包含 4 节课。
「面向智能体开发者的异步 Python」这节课中我会学到什么?
asyncio 基础、async def、await、事件循环——理解异步模型 你通过在浏览器中直接运行的动手代码来练习 AI Agents,全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 AI Agents 需要有经验吗?
无需任何先前经验。CoddyKit 上的 AI Agents 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 1 节课,共 4 节。
「面向智能体开发者的异步 Python」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 AI Agents 课中编写并运行代码吗?
能。每节 AI Agents 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。
此课程中的所有课时
- 面向智能体开发者的异步 Python
- 事件队列与消息代理
- 非阻塞并行工具执行
- 异步智能体框架:LangChain 及更多