那晚,我的爬虫又崩了
凌晨两点,我被报警短信震醒。线上爬虫服务又挂了。
打开监控面板,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?
讲清楚区别之前,得先明白这两个东西分别是什么。
await是一个等待操作符。
它的行为很简单:“停在这里,等到这个异步操作完成,把结果给我。”
data = await fetch(url) # 程序停在这里,fetch完成之前不往下走 print(data)
注意,await不会创建任何东西。它只是"等待"一个已经存在的异步操作。
打个比方:你打电话给客服,客服说"请稍等,我帮你查一下",然后你拿着电话等在原地——这就是await。你什么都没创建,只是在等。
那么问题来了:我写了tasks.append(fetch(url)),那不是创建了任务吗?
不是。fetch(url)返回的是一个协程对象(coroutine object),不是任务。协程对象就像一个"待办事项清单",上面写着"我要去请求这个URL"。但这个清单本身并不会被执行,除非你主动去驱动它。
你可以这样理解:
- 协程对象 = 一张写了计划的纸条
await= 站在原地,把纸条上的事情做完再走create_task= 把纸条交给一个助手,让他去做,你继续干别的
那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的陷阱
如果你以为掌握了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,什么时候用create_task?
总结下来,规则其实不复杂:
用await的场景:
- 你需要这个操作的结果才能继续往下走
- 单个异步操作,不需要并发
- 操作本身很快,不值得额外创建任务
比如读取数据库的一行记录、调用一个API取回配置、等待用户输入——这些都是"拿到结果再做事"的场景。
用create_task的场景:
- 你不需要立刻拿到结果,可以后面再取
- 需要并发执行多个异步操作
- 想让某些操作在后台运行,不阻塞主流程
比如批量请求多个API、同时下载多个文件、先发送通知但不等待反馈——这些都是"先发起再说,回头再取结果"的场景。
一个简单的判断标准: 写下await的时候问自己一句——“如果我不等这个结果,下面的代码能跑吗?”
- 能跑 → 用
create_task - 不能跑 → 用
await
那晚之后,我改成了这样
那天凌晨四点,我把代码改成了这样:
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语法"到"真的理解它在干什么",中间隔了无数次线上事故和凌晨排查。希望读完这篇文章的你,不用再走这些弯路。
下次写异步代码的时候,多问自己一句:这个操作,我是该等它,还是该派出去? 答案清楚了,代码也就对了。