引言
提到Python多任务处理,很多人第一反应就是"并发"、""性能"、"异步"。说实话,这个领域水挺深,初学者容易踩坑,资深开发者也会因为选型不当翻车。我见过太多人拿着多线程去跑CPU密集任务,结果发现性能不但没提升反而倒退;也见过有人在异步框架里塞了一个同步阻塞的数据库查询,整个事件循环直接卡死。这篇文章我想把Python多任务处理的几种主流方案——多线程、多进程、asyncio协程——放在一起讲透,包括底层原理、适用场景、实战代码、性能对比,以及我在实际项目里踩过的坑和排错思路。不管你是刚开始接触这个概念的新手,还是已经用过多任务处理但想深入理解其原理的开发者,这篇文章应该都能给你一些参考价值。
1. 多任务处理的核心困境:为什么Python这么特殊
1.1 从一场生产事故说起
去年有一个线上服务频繁出现响应超时,查了半天,最后定位到问题根源在一个很不起眼的环节:某个数据处理模块用同步循环去请求外部接口,一次要等两秒钟,十个数据就是二十秒,前端页面直接卡死。这个场景大家都遇到过吧。当时有人提议加多线程,有人提议换异步框架,还有人直接说把这段逻辑用其他语言重写。这其实折射出Python一个很大的特点:它能让你很轻松地写出多任务的代码,但写完之后效果如何,取决于你对这门语言底层机制的理解。
要说清楚这个问题,首先得理解一个基本事实:进程是资源分配的最小单位,线程是CPU调度的基本单位。但在Python的世界里,这两者的表现力并不一样,尤其是在面对CPU密集型和IO密集型任务的时候,差异会非常明显。严格来说,多任务处理解决的核心问题只有一个——让CPU在等待事情完成的时候不要闲着。至于"事情"到底是计算还是网络请求,决定了你应该选哪条路。
1.2 三种方案:多线程、多进程、异步IO,它们到底在忙什么
很多初学者分不清这三个概念的关系。我打个比方吧。
多线程就像一家餐厅雇了十个服务员,大家共用一个大厨房。客人点完菜,服务员去等菜端过来,这个等待过程你可能可以服务下一位客人,但真正做菜的大厨只有一位,而且所有服务员都得挤在这一个厨房里。
多进程则像是开了十家分店,每家有自己独立的厨房和厨师,互不干扰,但是每家店的装修、水电、食材库存全套都要独立配齐,成本高,协作也麻烦。
asyncio更像是用一个非常高效的服务员,一个人能同时盯几十张桌子的需求,他不用站在那里等每一道菜做完,而是记下每桌的需求,菜好了的时候再去端。如果你安排得当,他一个人就能顶十个普通服务员。
这个类比大致对应了技术层面的核心差异。多线程共享同一块内存空间,轻量但受GIL限制;多进程拥有独立的内存空间,能真正利用多核CPU,但通信成本高;asyncio在单线程内通过事件循环实现高并发,适合处理大量IO等待但几乎不消耗CPU尼的任务。
1.3 先认清任务性质:CPU密集型还是IO密集型
这个判断直接决定了你的技术选型,可以说比任何框架、任何代码技巧都重要。
CPU密集型任务的典型特征是:大量数学计算、图像处理、数据压缩、编码转换等等。这类任务占着CPU不停的算,几乎不等待外部资源,它的瓶颈是CPU计算能力本身。
IO密集型任务则刚好相反:网络请求、文件读写、数据库查询、用户输入等等。这类任务的特点是,大部分时间都在等待外部系统返回数据,CPU本身是闲着的,真正占用资源的是那些等待中的连接。
让我把两类任务的典型特征和对应方案整理成一个表:
| 任务类型 | 典型场景 | 瓶颈资源 | 推荐方案 | 原因 |
|---|---|---|---|---|
| CPU密集型 | 数值计算、图像滤镜、压缩 | CPU核心 | 多进程 | 可以绕过GIL,利用多核 |
| IO密集型 | 爬虫、API调用、文件读写 | 等待时间 | 多线程/协程 | 把等待时间用来处理其他任务 |
| 混合型 | 先算后等、先等后算 | 两者兼顾 | 多进程+协程组合 | 各取所长 |
这个判断不能靠感觉,我一般会先用一个简单的测试脚本跑一下。比如对一段计算代码,分别测一下单线程、多线程、多进程的耗时,差不多就能得出结论。很多时候,多线程跑CPU密集型的耗时反而比单线程更长,因为线程切换本身就有开销,再加上GIL的存在,两个线程交替执行同一个计算,就是纯粹的浪费时间。
2. 多线程:IO密集任务的经典选择
2.1 GIL到底是什么,它到底锁死了什么
说到Python多线程,GIL是一个绕不开的话题,网上相关的讨论已经很多了,但我觉得还是得从最本质的层面再说一遍,因为很多人的理解其实是有点偏差的。
GIL的完整称呼是全局解释器锁(Global Interpreter Lock),它的作用是保证同一时刻只有一个线程在解释器中执行字节码。注意,这里说的是"解释器中执行字节码",而不是"多线程完全无法并行"。这意味着:在纯计算场景,多线程确实无法利用多核;但在IO等待场景,线程在等待的时候会释放GIL,让别的线程去执行字节码,这就让多线程在IO密集型任务中有了用武之地。
这也就是为什么很多刚入门的人心里有个疑惑:既然GIL限制这么明显,为什么多线程还存在?答案就是,因为大多数真实业务场景都是IO密集型的。一次网络请求的等待时间可能是几十毫秒甚至几秒,而这几百毫秒里,如果只有单线程,CPU就是在白白等,多线程就能让CPU去处理别的请求。
不过我要提醒一点,GIL的锁定粒度是字节码级,但一个字节码操作里面可能包含多个Python操作。比如a += 1,它其实涉及读取、计算、赋值三个步骤,在某个时刻线程A可能刚读完a的值,还没来得及写回去,就被调度走了,线程B也读到同一个值,两边各自加一,最后结果只加了一次。这就是竞态条件,也是下面我要讲的锁的由来。
2.2 原生线程还是线程池:我建议直接用ThreadPoolExecutor
很早之前用Python写并发,大家习惯用threading.Thread手动创建线程,配合queue.Queue做任务分发,代码量不小,而且线程的管理、回收、异常处理都是坑。后来concurrent.futures.ThreadPoolExecutor出现之后,我基本就很少手写裸线程了。
线程池的核心优势不用多说,任务提交进去就不管了,线程自己复用,结果用Future对象取回。我以一个非常常见的场景为例,批量下载多个文件,来看一下这段代码怎么组织:
import concurrent.futures import urllib.request URLS = [ "https://example.com/file1.tar.gz", "https://example.com/file2.tar.gz", # ... 更多URL ] def download(url, save_path): with urllib.request.urlopen(url) as response: data = response.read() with open(save_path, "wb") as f: f.write(data) return url with concurrent.futures.ThreadPoolExecutor(max_workers=8) as executor: future_to_url = {executor.submit(download, url, f"{url.split('/')[-1]}"): url for url in URLS} for future in concurrent.futures.as_completed(future_to_url): url = future_to_url[future] try: result = future.result() print(f"下载完成: {result}") except Exception as exc: print(f"下载失败: {url}, 原因: {exc}")这段代码看似简单,里面的门道不少。as_completed会在任何一个Future完成时返回,这样主线程不会停滞在等待某一个慢请求上,整个循环会在所有任务完成后自然结束。future.result()必须放在try块里,因为子线程里抛出的异常不会直接传播到主线程,而是被封存在Future对象里,如果不显式调用result(),你可能根本不知道哪个任务失败了。
关于max_workers的取值,没有标准答案。我实践下来的经验是:如果你的任务主要是在等待网络响应,可以设成CPU核心数的5到10倍甚至更高,因为线程在等待时并不占用CPU;但如果任务里还夹杂着一些本地数据处理,这个值可以稍微保守一点。很多新手上来就设100个线程,结果发现目标服务器限流了,大量的连接被拒,反而拖慢了整体进度。
2.3 锁、信号量与线程安全:不写锁的多线程迟早翻车
说完线程池的便利,必须提一个不方便的地方:共享可变状态。如果你在多个线程里同时修改同一个列表、字典,或者访问同一个全局计数器,一定要加锁,或者使用线程安全的数据结构。
举个餐馆后厨的例子,你想想看,多个服务员都跑去同一个出菜口端菜,如果没有规则,几个人同时伸手,菜就撒了。加锁就是给这个出菜口上了一把小锁,一次只能让一个人拿。Python里的threading.Lock就是干这个的。
import threading counter = 0 lock = threading.Lock() def increment(): global counter for _ in range(100000): with lock: counter += 1 threads = [threading.Thread(target=increment) for _ in range(10)] for t in threads: t.start() for t in threads: t.join() print(counter)这里必须说一下,with lock:是上下文管理器写法,等价于lock.acquire()和lock.release()的组合。如果你不在异常路径里释放锁,很可能造成死锁,别的线程永远等下去。用with可以自动保证释放,这是我觉得最稳妥的写法。
另外,Python还提供了queue.Queue作为线程安全的队列,它的get()和put()方法内部已经有合适级别的安全机制。我建议生产者和消费者模式中优先使用队列来传递数据,而不是直接操作共享列表。本质上,用队列代替裸共享状态,是解决多线程数据竞争最优雅的方式之一。
3. 多进程:CPU密集型任务的王道
3.1 ProcessPoolExecutor与fork的安全问题
在Python的多进程编程中,有一个经典问题一直在被讨论:fork到底安不安全。简单解释一下,fork()是从操作系统层面复制当前进程的地址空间来创建一个子进程,父子进程一开始拥有一模一样的内存内容。对于单线程进程来说,fork是安全的;但如果父进程里已经启动了多个线程,这些线程在子进程里的状态会非常混乱,不同平台表现还不一样,这就导致了各种诡异的问题。
我自己的原则是:在多线程已经启动后,绝不在主进程中直接使用fork创建进程。在Python的multiprocessing模块里,有三种启动方式:fork(Linux默认)、spawn(Windows默认,Python 3.8+也算跨平台通用)、forkserver(介于两者之间)。实际开发中,我一般建议用spawn,它虽然启动慢一点,但每次都会重新加载一次Python解释器,不存在跨进程状态拷贝的问题,安全性最高。
接着说说ProcessPoolExecutor,它和多线程的ThreadPoolExecutor用起来几乎一模一样,只是底层机制不同。我这里有一个图像处理的实际例子,对一批图片做高斯模糊,然后批量保存:
from concurrent.futures import ProcessPoolExecutor from PIL import Image, ImageFilter import os def process_image(filename): with Image.open(filename) as img: filtered = img.filter(ImageFilter.GaussianBlur(radius=2)) output_name = "processed_" + os.path.basename(filename) filtered.save(output_name) return filename with ProcessPoolExecutor(max_workers=4) as executor: futures = [executor.submit(process_image, f) for f in os.listdir(".") if f.endswith(".jpg")] for future in concurrent.futures.as_completed(futures): try: result = future.result() print(f"处理完成: {result}") except Exception as exc: print(f"处理失败: {exc}")这里max_workers的选择比多线程更讲究,因为进程数跟CPU核心数是直接相关的。一般来说,设成CPU核心数就够了,再多反而会因为频繁切换造成额外开销。你可以在笔记本上用os.cpu_count()看看你的机器到底有几个核心。如果你有8个逻辑核心,用4个或8个都可以接受,具体可以小流量压一下测试。
3.2 进程间通信与序列化开销
多进程之间不共享内存数据,这让编程模型更安全,但也带来了新的成本——每次向子进程传参、从子进程收结果,都逃不掉序列化和反序列化。Python默认用pickle来做这个事,它对很多对象都支持,但也常常在自定义类上翻车。
之前踩过一个坑:把一个大DataFrame用ProcessPoolExecutor提交给子进程处理,每次要传600多MB的数据,导致性能不但没提升,反而因为序列化直接炸了内存。这个是我最想说的一点,多进程不是银弹,它的好用程度取决于任务边界是否清晰。如果任务之间的输入输出数据量特别大,或者每个子任务本身就很短,那么序列化的开销会远超并行带来的收益。这种情况下,更好的思路是用共享内存、文件或者消息队列(例如multiprocessing.Queue)来做数据流转,而不是简单地把数据塞进函数参数里。
下面是个简单的例子,展示如何通过队列在工作进程之间传数据:
import multiprocessing as mp def worker(q_in, q_out): while True: item = q_in.get() if item is None: q_out.put(None) break q_out.put(item * item) if __name__ == "__main__": q_in = mp.Queue() q_out = mp.Queue() p = mp.Process(target=worker, args=(q_in, q_out)) p.start() for i in range(100): q_in.put(i) q_in.put(None) q_out.get() # 拿到None,表示结束注意q_in.put(None)这个哨兵值的用法,它是终止子进程的常用技巧。否则子进程会一直阻塞在q_in.get()上,主进程结束时它也会变成孤儿进程。这种细节,实际项目中非常常见,社区里也有很多人因为忘了这个步骤导致进程悬挂。
3.3 进程池的资源管理:那次把所有CPU打满的教训
这里想分享一个实际踩坑的经历。有一段时间我需要批量处理大量PDF文件,提取文本内容,然后做关键词检索。一开始我图省事,ProcessPoolExecutor设的max_workers=16,想着反正内存大,多开几个进程跑得快。结果程序跑起来之后,全机器CPU直接打满,鼠标都在屏幕上飘移,内存也随着每个进程读取PDF文件开始暴涨,险些把整个开发机的其他服务全拖垮。
问题的本质在于:每创建一个进程,Python解释器本身占用的内存再加上任务数据的内存,是会成倍累积的。16个进程,哪怕每个进程只加载一个大文件进来,也很容易把物理内存耗尽。而且PDF文本提取本身有一部分计算在底层C库中执行,不受GIL保护,多进程的优势并没有想象中那么大。
那之后我的做法是:先拿一块小样本测试,统计单个任务的平均耗时和峰值内存,再根据这台机器的总内存倒推合适的进程数。比如机器有8核16G内存,单个任务峰值内存500MB,那进程数最多设4个,留出足够余量给系统和其他服务。这个原则后来被我写成了一个简单的公式:安全进程数 = min(CPU核心数, 可用内存 / 单任务峰值内存)。虽然听上去简单,但管用得很。
4. asyncio:单线程内的高并发魔术
4.1 事件循环、协程与await的本质
很多人一听asyncio就头大,觉得它跟多线程多进程完全不是一个世界的东西。其实只要理解了事件循环(Event Loop)的概念,asyncio就成功了一大半。
事件循环就像一个中央调度室,你可以往里面塞很多"任务",这个任务就是协程(coroutine)。当某个协程执行到一个IO等待操作时,关键一步是让出控制权,这个操作对应的写法就是await。事件循环收下这个请求后,会把它登记到等待列表里,继续去执行其他协程。当等待的IO完成之后,事件循环会回头唤起这个协程,让它从await的下一行继续运行。
用生活化的比喻就是:一个服务员不停地在各个餐桌之间转悠,看到谁举手了就过去记下需求,然后说"稍等,我去上菜",接着去另一桌。不用像同步阻塞思维那样,服务员站在一桌旁边等十分钟直到上菜完毕。
这个模型对于IO密集型任务特别高效,因为它几乎没有"额外线程切换"的开销。多线程的每一次切换都要操作系统参与,涉及上下文保存和恢复;而asyncio的切换是在用户态完成的,代价极小。
要正确编写asyncio代码,有几个认知需要建立起来。你写async def定义的函数不能直接调用,必须通过事件循环调度它;你await的必须是一个可等待对象(一般是另一个协程或者Future);事件循环是单线程的,所以在协程里绝对不能放置同步阻塞的任务,否则整个循环都会卡住。
4.2 一个异步爬虫示例,看asyncio如何编排并发任务
我来写一个相对完整的异步HTTP请求脚本,使用aiohttp库并发请求多个API,并把结果汇总。
import asyncio import aiohttp async def fetch(session, url): async with session.get(url) as resp: return await resp.json() async def main(): urls = [ "https://api.example.com/posts/1", "https://api.example.com/posts/2", "https://api.example.com/posts/3", ] async with aiohttp.ClientSession() as session: tasks = [asyncio.create_task(fetch(session, url)) for url in urls] results = await asyncio.gather(*tasks) for result in results: print(result["title"]) if __name__ == "__main__": asyncio.run(main())这里有几个细节值得注意。第一,aiohttp.ClientSession本身也是一个异步上下文管理器,它会维护一个连接池,所以必须在同一个事件循环里使用,不可以跨线程或跨循环复用。第二,asyncio.create_task是Python 3.7之后推荐写法,它会立刻把协程包裹成一个Task对象并安排到事件循环中,但注意,这只是"定了计划",还需要await gather(...)才能真正等它们全部完成。第三,asyncio.run()是Python 3.7之后的标准入口,帮你自动创建和关闭事件循环。
实际中,有些新手会把create_task和asyncio.sleep搭配做一些定时任务,需要注意底层逻辑:如果你在任务里await asyncio.sleep(0),相当于是主动让出当前协程的控制权,给事件循环一个机会去跑其他任务,这个技巧在排查"协程卡住"的问题时特别管用。
4.3 协程与多线程的对比:什么时候用哪个
很多人都有这个困惑:既然asyncio如此轻量高效,多线程是不是没有存在意义了?其实不是。关键区别在于:asyncio的单线程模型决定了它不适合处理CPU密集计算,哪怕你把计算封装成一个纯函数也一样,因为计算本身会占用事件循环,其他等待中的协程全都会被阻塞掉,那就得不偿失了。
我把两者的对比简单总结一下:
| 维度 | 多线程 | asyncio |
|---|---|---|
| 并发模型 | 抢占式,由操作系统调度 | 协作式,由事件循环调度 |
| 切换开销 | 系统级,较大 | 用户级,极小 |
| 数据共享 | 需要锁,存在竞态 | 单线程天然共享,不用锁 |
| 适合任务 | IO密集型,少量计算 | IO密集型,无CPU计算 |
| 代码风格 | 函数式,相对直观 | 异步风格,需要适应 |
| 现场运维 | 线程栈难以定位 | 有调试模式的挂钩 |
我在实际项目里两者都用,但有一条很明确的原则:只要涉及的第三方库没有提供异步版本,或者你不确定它内部是否阻塞,优先考虑用线程池把它包起来,放到asyncio之外单独执行,而不是硬塞进协程里await。
比如你有一个老旧的数据库驱动,只提供同步接口,你把一个查询直接放进协程里,如果数据库响应慢了,事件循环就卡住了。此时正确做法是把数据库查询提交给asyncio.to_thread(Python 3.9+)或loop.run_in_executor,让它在独立线程中执行,执行完成后回调回事件循环。
5. 混合架构:我如何在一套系统里同时用三种方案
5.1 一个典型的数据处理系统的架构拆解
上面分别讲了三种方案的原理和用法,但真实世界里的系统极少只用其中一种。绝大多数情况下都是组合使用,各取所长。我这里有一个实际项目做参考:某跨平台系统中的数据采集模块,它需要从几百个外部数据源定时抓取数据,抓取之后还需要做解析、去重、聚合计算,然后写回数据库。
这套系统的核心设计是这样的:
第一层,asyncio负责高并发网络请求。因为数据源数量多,每个请求的等待时间差异很大,有的几十毫秒,有的几十秒,如果用同步顺序请求,一轮采集可能要跑一个小时。asyncio在这里可以轻松挂起上千个并发连接,只要事件循环不被阻塞,整体采集速度会提升一个数量级。
第二层,解析和聚合计算交给多进程。因为从各个源抓回来的原始数据可能是JSON、XML、CSV格式,解析之后还需要做一些正则清洗、编码转换、字段映射,这些计算虽然不是特别重,但量大了也会耗时。把它放到进程池里,让多个CPU核心并行计算,能明显压缩整体耗时。
第三层,落库和日志写入使用多线程。数据库连接往往是同步的,而且同一时刻数据库能承受的连接数有限,如果放在asyncio里会卡住事件循环,如果放在进程里又没有必要建那么多独立进程。用线程池固定几个工作线程来写库,排队慢慢写入即可。
5.2 组合架构的代码骨架与关键衔接点
下面我写一个简化版本,展示asyncio如何与线程池、进程池配合工作:
import asyncio import aiohttp from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor async def fetch_all_urls(urls): """用asyncio并发抓取所有URL""" async with aiohttp.ClientSession() as session: async def fetch_one(url): async with session.get(url) as resp: return await resp.text() results = await asyncio.gather(*[fetch_one(url) for url in urls]) return results def parse_batch(data): """解析数据,放在进程池里跑""" # 略,应为CPU密集型的解析逻辑 return [len(x) for x in data] async def main(): urls = [f"https://example.com/data/{i}" for i in range(100)] raw_data = await fetch_all_urls(urls) loop = asyncio.get_running_loop() # 用进程池做CPU密集解析 with ProcessPoolExecutor(max_workers=4) as pool: parsed = await loop.run_in_executor(pool, parse_batch, raw_data) # 用线程池做数据库写入 with ThreadPoolExecutor(max_workers=2) as pool: for item in parsed: await loop.run_in_executor(pool, write_to_db, item) def write_to_db(item): # 略,同步的数据库写入 pass if __name__ == "__main__": asyncio.run(main())这个骨架的价值在于:网络IO密集部分用了asyncio,CPU密集解析部分用了进程池,同步数据库写入用了线程池,每个组件都处于自己最擅长的位置。
不过有一点我踩过坑之后想提醒大家:asyncio里的loop.run_in_executor每次调用都会向线程池或进程池提交新任务,但同一时刻允许并发的任务数受线程池/进程池大小限制。如果你在循环里提交了1万个任务,底层池子只有8个worker,那其余任务会排队等待,但这种排队是在事件循环内部自动管理的,不会阻塞主循环,只是总耗时可能会比预期更长。做一个粗略评估:整个任务的吞吐量不应该超过最慢的那一层处理能力。
5.3 性能测试:别用time.time(),要用正确的工具
聊完架构,聊一个非常重要但经常被忽略的话题——如何测试多任务处理之后的性能提升。很多人图省事,直接在代码前后包两个时间戳,time.time()一减就完事。我建议不要这样干,至少有这么多问题它测不出来:上下文切换开销、锁竞争、进程创建耗时、IO等待分布。
我一般的做法是:写一个独立脚本,模拟实际任务(包括IO等待和计算部分),分别跑多个方案,使用time.perf_counter()进行基准测试,每个方案重复多轮取平均值。在测量多进程性能时,一定要排除启动进程的时间,否则你会误以为多进程很慢。
此外Python自带的cProfile可以查看每个函数的调用次数和耗时,futures的as_completed可以对任务的完成时间分布做个可视化统计。用faulthandler能帮助排查死锁和卡死问题,用py-spy可以dump出当前进程的调用栈,看到卡在哪一行,这个工具在线上排查多线程问题时几乎救了我的命。
6. 常见问题与排查技巧实录
6.1 我遇到过的高频问题速查表
下面这个表记录的是我自己在不同项目里踩过的坑,做成一张问题对照表,方便大家快速定位:
| 现象 | 可能原因 | 排查方式 | 解决方案 |
|---|---|---|---|
| 多线程比单线程还慢 | 任务其实是CPU密集型 | 用py-spy看现场,测CPU占用 | 换多进程;如果是计算量大,考虑底层用C扩展 |
| 多进程运行时内存暴涨 | 每个进程复制了大量数据 | 检查每个进程的RSS内存 | 减少进程数,或用共享内存/队列传小数据 |
| 协程全部卡住不动 | 协程里混入了同步阻塞调用 | 访问协程卡死的堆栈 | 使用asyncio.to_thread把阻塞调用挪走 |
| 多线程写同一个文件内容乱掉 | 缺少锁或使用了非原子写 | 检查写文件逻辑 | 加锁,或改用队列单线程写 |
| 程序退出时卡在等待 | 有任务没被标记完成 | 使用wait而不是gather,并设置超时 | 显式设置asyncio.wait_for超时 |
| 子进程崩溃且无输出 | 子进程里抛了异常,主进程忽略了 | 在回调里加上future.result()的异常捕获 | 给进程池的任务包装一层带日志的函数 |
| 开线程过多导致系统无法响应 | 无上限创建线程 | 观察线程数统计 | 一定用线程池,控制max_workers |
这张表不是一次写出来的,很多条目都是我改了N轮之后才沉淀出来的经验。排查这类问题最忌讳的是靠猜,用工具看调用栈永远是第一要点。
6.2 用py-spy定位卡死问题的实战记录
我记得有一次线上服务出现"假死"状态,网页能打开但所有接口都转圈。我们用py-spy dump --pid <进程ID>去抓取主进程的调用栈时发现,几乎所有的线程都阻塞在同一个queue.Queue().get()上面,队列里没有新数据,而且队列的生产者线程早就因为异常退出了。
这个问题的本质是:消费者线程在queue.get()上无限期等待,但生产者异常退出后不会再放任何消息进来,所有等待线程就成了"不死不活的僵尸线程"。解决方案很简单,生产者退出前必须往队列里放一个None哨兵值,消费者拿到None就退出循环。另一种方案是使用queue.get(timeout=5),等五秒拿不到就重新检查退出标志,这样线程就不会无限期卡住。
排查多任务程序的卡死问题,还有一个非常高效的思路:给每一个长任务加上超时机制。具体到Python,concurrent.futures里可以用future.result(timeout=10)来限时等待结果,asyncio里可以用asyncio.wait_for(coro, timeout=10)。超时后你可以选择放弃这个任务、记录日志、或者重新调度,而不是让整个程序永远卡在那里。
6.3 几个让我印象深刻的经验教训
说了这么多,最后聊几点个人感受比较深的东西。
第一,不要过度并发。很多开发者的本能是一看到性能瓶颈就必须加线程加进程,但我的经验是,并发规模的临界值往往取决于下游依赖的承受能力,而不是你本地机器的配置。在爬虫场景里,目标服务器的连接数限制、频率限制才是天花板;在数据库场景里,数据库连接池大小才是瓶颈。一上来就把max_workers拉到几十上百,很容易触发对方限流甚至封禁,还没用更多线程,速度反而更慢了。反过来,用协程做大量网络请求也要注意连接数控制,比如aiohttp.TCPConnector(limit=100)限制同时打开的连接数,给双方都留一点缓冲空间。
第二,多进程虽然能绕开GIL,但它不是万能的,序列化开销和数据传递开销都很容易成为新的瓶颈。如果你的子任务很小而数据量很大,多进程的代价甚至可能超过收益。我在项目里有一条约定:只有在单个任务运行时间超过1秒,并且输入输出数据小于几MB的时候,才优先选择多进程;否则用多线程或者asyncio。
第三,多任务的调试体验一定比单线程差,这是客观现实。线程崩溃不会直接让你看到错误栈,进程崩溃也常常悄无声息,协程里出现的异常如果不是被await碰到,可能永远不会被抛出。所以写多任务代码时,第一要务不是追求性能,而是先保证异常可见性。我的做法是:每个被提交的函数都包一层try-except,把异常、任务参数、线程/进程信息都记录到结构化日志里。虽然代码显得啰嗦了一点,但在线上排查问题的时候能省下大量的猜谜时间。
第四,不建议一开始就用最复杂的技术方案。如果你的业务只有几十个请求,同步写就够用了;如果是几百个请求,线程池通常就能解决;只有上千个、上万个并发请求,才是asyncio发挥价值的时候。多任务处理是一种手段,不是目的,盲目地上高并发架构,只会让自己陷入调试地狱。这一点是我在项目里测试了很多轮之后最深的体会。