
凌晨两点,我被报警短信震醒。线上爬虫服务又挂了。
打开监控面板,CPU跑满,内存飙升,连接数异常。重启之后不到十分钟,再次崩溃。翻日志只看到一行:TimeoutError,没有更多信息。
这已经是一个月内的第三次了。
代码我看过无数遍,逻辑清晰,结构合理。用了asyncio,用了aiohttp,该await的地方都await了,按理说不该出问题。
直到我盯着那几行核心代码看了整整两个小时,才突然意识到一个被我反复忽略的事实:
await和asyncio.create_task,根本就不是一回事。而我,一直在把它们混着用。
先给你看看我当时写的代码(简化版):
import asyncio
import aiohttp
async def fetch(url):
async with aiohttp.ClientSession() as session:
async with session.get(url) as response:
return await response.text()
async def main():
urls = ['https://api.example.com/data'] * 100
for url in urls:
result = await fetch(url) # 一个个来?
await save_to_db(result) # 一个个存?
asyncio.run(main())我看到这代码的第一反应是:这跟同步有什么区别?
于是我"优化"了一下:
async def main():
urls = ['https://api.example.com/data'] * 100
tasks = []
for url in urls:
tasks.append(fetch(url))
for task in tasks:
result = await task # 还是一个个等?
await save_to_db(result)看起来创建了一堆任务对吧?但实际上,这里连一个并发都没有。
为什么?因为你虽然把协程对象放进了列表,但在await task的那一刻,程序还是老老实实等它执行完,再执行下一个await。整个流程依然是串行的。
这就好比你点了十份外卖,但每一份都等送到吃完才点下一份——那跟只点一份有什么区别?
当时的我还自我安慰说"用了异步应该比同步快吧",结果压测数据打脸打得啪啪响。QPS跟同步代码几乎一样,该超时的照样超时。
讲清楚区别之前,得先明白这两个东西分别是什么。
await是一个等待操作符。
它的行为很简单:“停在这里,等到这个异步操作完成,把结果给我。”
data = await fetch(url)
# 程序停在这里,fetch完成之前不往下走
print(data)注意,await不会创建任何东西。它只是"等待"一个已经存在的异步操作。
打个比方:你打电话给客服,客服说"请稍等,我帮你查一下",然后你拿着电话等在原地——这就是await。你什么都没创建,只是在等。
那么问题来了:我写了tasks.append(fetch(url)),那不是创建了任务吗?
不是。fetch(url)返回的是一个协程对象(coroutine object),不是任务。协程对象就像一个"待办事项清单",上面写着"我要去请求这个URL"。但这个清单本身并不会被执行,除非你主动去驱动它。
你可以这样理解:
await = 站在原地,把纸条上的事情做完再走create_task = 把纸条交给一个助手,让他去做,你继续干别的asyncio.create_task()才是真正"创建任务"的东西。
它会把你给的协程包装成一个Task对象,然后立即调度到事件循环中执行。注意"立即"这两个字——任务从创建那一刻起就开始跑了,不需要你手动await。
task = asyncio.create_task(fetch(url))
# 到这里,fetch(url)已经开始执行了
# 程序不会阻塞,立刻往下走回到打电话的比喻:create_task不是让你拿着电话等,而是你把事情交代给助手,助手马上开始处理,你放下电话该干嘛干嘛。
关键区别来了:
操作 | 是否阻塞 | 是否创建任务 | 是否立即执行 |
|---|---|---|---|
await 协程 | 是(等结果) | 否 | 立刻执行并等待 |
create_task(协程) | 否 | 是 | 立刻执行但不等待 |
用create_task改一下前面的爬虫:
async def main():
urls = ['https://api.example.com/data'] * 100
tasks = []
for url in urls:
task = asyncio.create_task(fetch(url))
tasks.append(task)
# 到这里,所有请求已经同时发出了
# 现在可以等它们全部完成
results = await asyncio.gather(*tasks)
for result in results:
await save_to_db(result)这才叫真正的并发。
100个请求同时发出,等待时间从"100次网络往返"变成了"1次网络往返中最慢的那一次"。在我的测试环境里,响应时间从45秒直接降到了2.3秒。
如果你以为掌握了create_task就万事大吉,那就太天真了。我踩过的坑远不止这一个。
async def main():
asyncio.create_task(fetch(url)) # 创建了任务
# 但程序立刻结束了,任务根本没跑完
asyncio.run(main())asyncio.run()在main()结束后会关闭事件循环。如果还有任务没完成,它们会被直接丢弃——连报错都没有。
你以为任务在跑,实际上它已经被悄悄杀掉了。
解决方案:用asyncio.gather()或者await task来确保任务完成。
async def main():
task = asyncio.create_task(fetch(url))
await task # 确保任务完成或者:
async def main():
tasks = [asyncio.create_task(fetch(url)) for _ in range(10)]
await asyncio.gather(*tasks) # 全部等待create_task创建的任务如果抛出了异常,而这个异常没有被捕获,它不会像普通函数那样直接报错让你看见。异常会保存在Task对象内部,等你await task或者task.result()的时候才会抛出来。
如果你不等待任务,异常就永远消失了。
async def broken():
raise ValueError("出错了")
async def main():
task = asyncio.create_task(broken())
# 这里啥也没发生,异常被吞了
await asyncio.sleep(1) # 程序安静地运行
asyncio.run(main()) # 没有任何报错这个问题非常隐蔽。我当时那个凌晨崩溃的爬虫,就是因为某个任务抛了异常没被捕获,然后任务一直卡在那里占着资源,最后把整个事件循环拖垮了。
正确做法是给每个任务加异常处理:
async def safe_fetch(url):
try:
return await fetch(url)
except Exception as e:
log.error(f"请求{url}失败: {e}")
return None
# 或者在获取结果时捕获
results = await asyncio.gather(*tasks, return_exceptions=True)
for result in results:
if isinstance(result, Exception):
log.error(f"任务失败: {result}")
else:
process(result)create_task的调度成本很低,但不是零。如果你一下子创建一万个任务,事件循环的调度器会忙不过来。每个任务都要维护状态、切换上下文,开销不可忽略。
我的爬虫曾经一次性创建了5000个任务去请求不同的API,结果事件循环的切换开销比网络请求本身还大。
解决方案是控制并发数量——用asyncio.Semaphore或者asyncio.gather分批处理。
sem = asyncio.Semaphore(100) # 最多同时100个
async def limited_fetch(url):
async with sem:
return await fetch(url)
tasks = [limited_fetch(url) for url in urls]
results = await asyncio.gather(*tasks)总结下来,规则其实不复杂:
用await的场景:
比如读取数据库的一行记录、调用一个API取回配置、等待用户输入——这些都是"拿到结果再做事"的场景。
用create_task的场景:
比如批量请求多个API、同时下载多个文件、先发送通知但不等待反馈——这些都是"先发起再说,回头再取结果"的场景。
一个简单的判断标准: 写下await的时候问自己一句——“如果我不等这个结果,下面的代码能跑吗?”
create_taskawait那天凌晨四点,我把代码改成了这样:
async def main():
urls = get_urls()
# 用Semaphore控制并发数
sem = asyncio.Semaphore(50)
async def fetch_with_limit(url):
async with sem:
try:
return await fetch(url)
except Exception as e:
log.error(f"URL失败 {url}: {e}")
return None
# 创建所有任务,立即开始执行
tasks = [asyncio.create_task(fetch_with_limit(url)) for url in urls]
# 等待所有任务完成,异常不抛出,而是作为结果返回
results = await asyncio.gather(*tasks, return_exceptions=True)
# 处理结果,把异常的单独拎出来
success = []
failed = []
for url, result in zip(urls, results):
if isinstance(result, Exception) or result is None:
failed.append(url)
else:
success.append(result)
log.info(f"成功{len(success)}条,失败{len(failed)}条")
# 失败的重试
if failed:
retry_failed(failed)
asyncio.run(main())改完之后跑了三天,再也没崩过。
除了上面说的,还有几个小坑顺带提一下,都是血泪教训:
1. 不要在循环里直接await异步函数。 这等于把并发变成了串行。要么先收集所有协程对象再gather,要么用create_task创建任务列表。
2. asyncio.run()和get_event_loop()别混用。 asyncio.run()每次会创建新的事件循环,如果你在里面用了get_event_loop()拿到的是不同的循环,任务可能调度不进去。
3. 异步函数里别用time.sleep(),用asyncio.sleep()。 time.sleep()会阻塞整个线程,事件循环直接卡死。当年我因为这个bug排查了一整天。
4. gather和wait的区别。 gather返回所有结果,wait返回完成/未完成的任务集合。简单场景用gather就够了,需要精细控制超时和取消时用wait。
await是"等结果",create_task是"派任务"。
前者让你停下脚步,后者让你同时做多件事。异步编程的核心就是在这两者之间找到平衡——该等的时候等,该并发的时候并发。
三年时间,我从"会写async/await语法"到"真的理解它在干什么",中间隔了无数次线上事故和凌晨排查。希望读完这篇文章的你,不用再走这些弯路。
下次写异步代码的时候,多问自己一句:这个操作,我是该等它,还是该派出去? 答案清楚了,代码也就对了。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。