ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Python并发编程全解析:进程、线程、协程与IO模型选型实战

Python并发编程全解析:进程、线程、协程与IO模型选型实战 提到Python并发编程很多人第一反应是“多线程让爬虫更快”或者“协程比线程更轻”再不然就是“进程池到底该开几个”。我在技术社区里看过太多类似的问题进程等待wait到底等的是什么线程死锁怎么定位Windows下某个进程被锁定重启服务就好了这些看似零散的问题背后其实都指向同一套底层逻辑——进程、线程、协程、IO模型之间的关系没理顺。这篇不打算按教科书顺序给你背概念而是用我这些年写Python并发工具、排查线上事故的真实经验把Python下的并发家族完整捋一遍。顺便解决三个最头疼的问题模型怎么选、参数怎么调、出问题怎么查。无论你是刚接触Python的新手还是被GIL、事件循环、进程池折磨过的老手后面这些内容应该都能直接拿去做参考。1. 先把并发家族认全进程、线程、协程到底在争什么1.1 进程资源隔离的“重装系统”进程是操作系统分配资源的基本单位。每个进程有自己独立的内存空间、文件描述符、环境变量甚至独立的解释器状态。在Python里用multiprocessing.Process创建子进程本质上等于启动了一个独立的Python解释器实例子进程和父进程之间不能直接共享普通变量。进程最大的优点是隔离性好一个子进程崩溃了不会直接把父进程拖垮子进程里拿到了什么脏数据也不会污染主进程的全局状态。缺点也明显——创建和切换进程的开销比线程高出一个数量级因为操作系统要重新分配地址空间、初始化运行时环境。我常用一个生活化类比进程就像“重装系统”。你要跑一个新进程操作系统就得重新配一套运行环境而线程则像是在已经开机的系统里再开一个窗口轻得多但也共享同一个桌面的状态。也正因为每个进程有独立的Python解释器multiprocessing可以绕开GIL的限制。多进程代码在CPU密集计算场景下是真的能利用多核的。不过代价是进程间通信IPC变复杂了不能直接读变量必须走队列、管道、共享内存这些机制。这点后面专门展开说。1.2 线程共享内存的“轻装上阵”线程是进程内的执行单元。同一个进程下的多个线程共享堆内存、全局变量和文件描述符创建线程的开销比进程小很多切换也更快。但Python里有个绕不开的坎GIL全局解释器锁。GIL保证同一时刻只有一个线程能执行Python字节码。注意它不是把整个解释器所有操作都锁住而是线程每执行一段字节码或经过一个时间片就会释放GIL让其他线程获取。所以如果一个线程正在等网络响应、等磁盘读写它会把GIL释放出来其他线程就能趁机执行。这是Python多线程在IO密集型场景下依然有价值的核心原因。我理解的GIL像一个“话筒”办公室里很多人可以同时办公但只有一个话筒一次只能一个人说话。如果有人在等电话回拨话筒就能交给别人用。这里有个常见的误区认为多线程一定比单线程快。对于CPU密集的纯计算任务threading不仅不会提速反而会因为上下文切换和锁竞争导致更慢。这时候应该考虑多进程或者直接优化算法本身而不是盲目上线程。1.3 协程单线程里的“时间片魔术师”协程不是操作系统实体它跑在单线程内部由语言层的调度器配合事件循环完成切换。Python的asyncio会给每个await点做一次主动让权遇到IO等待就先跑去执行其他任务等IO就绪了再回来继续。因为协程的切换不需要操作系统介入创建成本极低单进程里开几十万个协程都没什么压力线程开到几万个系统一般就开始告警了。协程的类比是一个特别擅长统筹的大厨锅里的菜还没熟他不会干等而是去切配菜、准备调料。但协程不是真正的并行只是并发。如果代码里有一个耗时的同步阻塞操作比如直接用requests.get()整个事件循环会被卡住因为阻塞操作不会主动让权。这是新手最容易踩的坑用了asyncio.sleep没问题换成了同步IO就全卡住。解决办法是把这个阻塞调用扔到线程池里执行或者改用原生的异步库。2. 进程篇multiprocessing 实战与 IPC 细节2.1 进程创建与生命周期管理multiprocessing.Process的用法不复杂但有几个坑我几乎每次带新人都会强调。第一在Windows上进程的创建方式是spawn也就是新进程会重新import一遍当前模块。所以必须把进程创建逻辑放在if __name__ __main__:里保护否则会无限递归创建子进程。Linux/macOS默认是fork的平台上没有这个问题但为了跨平台兼容代码里最好都加上这段保护。第二join()必须有。它对应热搜词里的“进程等待wait”作用是让父进程阻塞等待子进程结束并回收其退出状态。如果忘了join()主程序可能已经退出子进程还没跑完你就会发现日志丢失、结果没写进去。第三join()可以设置超时。from multiprocessing import Process import time def worker(name): print(f子进程 {name} 开始工作) time.sleep(2) print(f子进程 {name} 完成) if __name__ __main__: procs [] for i in range(3): p Process(targetworker, args(i,), namefworker-{i}) procs.append(p) p.start() for p in procs: p.join(timeout5) # 检查是否还有子进程没跑完 for p in procs: if p.is_alive(): print(f{p.name} 超时仍然存活需要终止) p.terminate() print(所有子进程处理完毕)注意join(timeout5)超时返回并不代表子进程结束。如果你后面的逻辑依赖子进程的结果最好在join之后检查is_alive()必要时terminate()。还有一个细节Process.daemon True的子进程会随主进程退出被强制终止适合后台日志清扫任务但不要在里面写重要数据否则可能来不及落盘就没了。2.2 进程间通信Queue、Pipe 与共享内存进程间通信IPC是很多人被卡住的地方。multiprocessing.Queue本质是一个线程安全和进程安全的FIFO队列底层用管道加锁实现。用Queue时放进去的对象必须能被pickle序列化。我踩过一个很经典的坑把一个定义在if __name__ __main__内部的匿名类实例放进Queue结果子进程拿到的数据反序列化失败。所以放入队列的对象最好都定义在模块顶层或者直接传基本数据类型。看一个生产者和消费者的标准写法from multiprocessing import Process, Queue def producer(q): for i in range(5): q.put(f消息 {i}) def consumer(q): while True: item q.get() if item is None: # 哨兵约定通知退出 break print(f收到{item}) if __name__ __main__: q Queue() p1 Process(targetproducer, args(q,)) p2 Process(targetconsumer, args(q,)) p1.start() p2.start() p1.join() q.put(None) # 发送停止信号 p2.join()这里有一个我在生产里遇到的教训producer和consumer如果用同一个Queue一定要约定好结束方式。很多人的程序卡死是因为消费者在q.get()上一直等待而生产者已经退出了队列里永远不会有新数据程序就像挂了一样。解决办法要么是传哨兵值要么用get(timeout...)。如果想在进程间传大量数据Queue可能不是最优解。因为每个对象都要序列化和反序列化大量数据时开销很大。此时可以考虑Pipe适合两个进程之间点对点通信或multiprocessing.shared_memory适合真正的大数据共享。共享内存虽然快但要用锁保护否则并发读写很容易出现逻辑错误。2.3 进程池与执行策略选择进程的创建和销毁开销很大所以需要进程池复用进程。multiprocessing.Pool是经典方案concurrent.futures.ProcessPoolExecutor则提供了更统一的接口。我一般推荐后者因为它和ThreadPoolExecutor的API一致用as_completed可以方便地处理任务完成顺序。from concurrent.futures import ProcessPoolExecutor, as_completed def square(n): return n * n if __name__ __main__: with ProcessPoolExecutor(max_workers4) as executor: futures [executor.submit(square, i) for i in range(100)] for future in as_completed(futures): print(future.result())进程池大小怎么定我的经验公式是CPU密集型任务max_workers设为CPU核心数或核心数1IO密集型任务用进程池其实是浪费因为IO等待不需要CPU进程反而比线程重。如果任务同时混合了计算和IO建议拆分计算部分扔进程池IO部分用协程或线程池再通过队列组合起来。还有一点要提醒在Jupyter或交互式环境里跑ProcessPoolExecutor很容易报“Function missing from__main__”。因为spawn模式下子进程需要重新导入任务函数而交互式环境里的函数无法被导入。解决方法很简单把任务函数写到独立.py文件里再从脚本调用。3. 线程篇GIL 阴影下的多线程真相3.1 threading 基础与锁的哲学threading.Thread创建线程很简单难的是共享数据的管理。线程之间共享全局变量这既是优势也是灾难源头。经典问题就是计数器丢失更新import threading counter 0 lock threading.Lock() def increment(): global counter for _ in range(10000): with lock: counter 1 threads [threading.Thread(targetincrement) for _ in range(10)] for t in threads: t.start() for t in threads: t.join() print(counter)如果去掉Lock最终counter大概率不是100000因为counter 1不是原子操作而是“读取-加一-写回”三步线程切换刚好发生在读取和写回之间时更新就丢失了。关于锁的哲学我自己的原则是锁粒度越小越好但不要小到逻辑不完整。用with lock:而不是手动acquire/release可以避免异常时忘记释放锁。Python里还有几个常用同步原语threading.Event一个线程发信号另一个线程等待信号适合做启动通知。threading.Semaphore限制同时执行某个操作的线程数量比如限制同时访问数据库的连接数。threading.RLock可重入锁同一线程可以多次acquire适合递归场景。如果一个函数内部调用了另一个也需要加锁的函数用普通Lock就会死锁RLock允许同一线程重复获取可以避免这个尴尬局面。3.2 线程池的合理配置与阻塞队列选择并发场景里ThreadPoolExecutor比手动建线程更好用因为线程池复用线程减少创建销毁的开销。线程数怎么定网上有“2 * CPU核心数 1”的经验公式但它只对IO密集任务有一定参考意义。我实际做法是先做压测从一个较小的线程数开始逐步增大观察吞吐量和延迟的曲线。吞吐量不再增长甚至回落就是最优线程数。盲目加大线程数只会增加上下文切换和锁竞争。一个容易被忽略的细节ThreadPoolExecutor内部的任务队列是无界队列。如果提交任务的速度远大于处理速度大量待执行任务会堆积在内存里最终把程序撑爆。Python的标准库没有暴露自定义队列大小的选项所以需要在提交端做准入控制。我的方案是配合threading.Semaphore限制在途任务数量import threading from concurrent.futures import ThreadPoolExecutor sem threading.Semaphore(20) def process(item): with sem: # 实际处理逻辑 pass这样即使外部疯狂提交任务真正进入执行状态的也始终有上限避免无界堆积。3.3 死锁与调试如何定位一个卡死的线程死锁的典型场景是线程A持有锁1等待锁2线程B持有锁2等待锁1。两个线程互相等待永远无法推进。Python里没有直接强制终止一个线程的API所以死锁发生后通常只能整体重启进程。这就意味着定位死锁比解决死锁更重要。我的排查工具顺序是先用faulthandler.dump_traceback_later(30)让程序在运行30秒后自动打印所有线程的调用栈。如果能看到两个线程分别卡在锁上死锁基本实锤。用py-spy dump --pid pid这是一个外部工具不需要侵入代码生产中也可以直接用。或者直接在代码里捕获当前帧import sys import threading import traceback def dump_thread_stack(): for thread_id, frame in sys._current_frames().items(): name threading.current_thread().name print(f线程 {name} (id{thread_id}):) traceback.print_stack(frame)看到一个线程在Lock.acquire处等待另一个线程也卡在获取另一把锁基本就能拼出死锁关系图。这时候不要只盯着代码改而是审视一下锁的获取顺序是否一致。如果所有线程都按照“先锁A再锁B”的顺序获取锁死锁就不会出现。另一种彻底避免死锁的方案是threading.Lock配合acquire(timeout...)超时获取失败就走备用逻辑宁可重试也不要无限等待。4. 协程篇asyncio 与事件循环4.1 从生成器到 async/awaitasyncio从Python 3.4开始引入最初用yield from生成器实现协程代码可读性很差。Python 3.5引入async/await语法后协程写起来终于像普通的顺序代码。import asyncio async def main(): print(开始) await asyncio.sleep(1) print(结束) asyncio.run(main())asyncio.run()是3.7之后官方推荐的入口它负责创建事件循环、运行协程、关闭事件循环。注意一个限制如果你正在一个运行中的事件循环里不能再调用asyncio.run()会直接抛RuntimeError。这时候应该用asyncio.create_task()创建任务或者await直接调用。协程和普通函数的区别在于调用main()不会执行函数体只会创建协程对象必须放进事件循环里运行才有效。这个区别很多新手都会迷惑看到“协程没有等待”的警告就很慌。4.2 协程间的并发控制与任务编排真正的并发场景要用asyncio.gather()把多个协程打包并发执行。比如模拟并发爬取十个URLimport asyncio async def crawl(url): # 模拟网络请求 await asyncio.sleep(1) return url async def main(): urls [fhttps://example.com/page/{i} for i in range(10)] results await asyncio.gather(*[crawl(u) for u in urls]) print(results) asyncio.run(main())这段代码总共耗时约1秒而不是10秒因为10个await asyncio.sleep(1)是并发执行的。gather()有个行为要知道如果某些任务返回异常gather会默认向外抛异常其他任务也会被取消。如果想让单个任务的异常不拖累整体可以改用asyncio.wait()或者给每个协程包一层try/except。突发高并发还需要限制速率我用asyncio.Semaphore做并发限流sem asyncio.Semaphore(10) async def safe_crawl(url): async with sem: return await crawl(url)这样即使一次提交1000个URL同时发起的网络请求也只有10个不容易把目标服务打挂。4.3 协程与线程的混合使用场景协程不是万能的。如果代码里有大量同步阻塞库的调用比如同步的数据库驱动、requests、文件读写事件循环就会被卡住。此时最好让这些阻塞操作在线程池里跑然后以协程方式等待结果。import asyncio import requests async def fetch_with_thread(url): response await asyncio.to_thread(requests.get, url) return response.textasyncio.to_thread是Python 3.9的工具底层就是用默认线程池执行函数。也可以用loop.run_in_executor()但to_thread语法更简洁。还有一种混合架构主事件循环负责高并发的网络连接需要CPU密集计算的子任务扔给ProcessPoolExecutor。这里有个坑进程池任务是阻塞的提交大量任务时事件循环可能会等得不耐烦。建议用独立线程专门管理和进程池交互避免阻塞主循环。5. IO 模型篇阻塞、非阻塞、多路复用与异步5.1 五种 IO 模型对比并发编程的本质其实是等待IO。网络请求、磁盘读写、数据库连接绝大部分时间都在等。理解IO模型比死记进程线程区别更能解决实际问题。IO模型基本行为是否阻塞Python对应阻塞IO等待数据就绪期间线程挂起是默认socket、普通文件读写非阻塞IO数据未就绪立即返回否socket.setblocking(False)IO多路复用同时监听多个连接批量等待就绪否select、poll、epoll、selectors信号驱动IO内核通过信号通知数据就绪否少用异步IO提交IO请求后完成全部操作后再通知否asyncio、IOCP、io_uring阻塞IO最简单但并发能力弱。每开一个阻塞连接至少要一个线程去等它线程一多系统就累了。非阻塞IO虽然不阻塞但需要你自己轮询成本也很高。多路复用是把“我该读哪个连接”的判断交给内核显著提升并发连接数上限。异步IO更进一步连数据的读写都由内核完成IO完成后再回调。Python的asyncio底层其实是“多路复用 协程调度”不是真正的内核异步IO但业务上已经很接近异步的效果了。5.2 多路复用与 select/epoll 在 Python 中的体现Python标准库里的selectors模块封装了不同平台的IO多路复用机制Linux上默认用epollWindows上用select或IOCP。asyncio的事件循环就构建在这层封装之上。如果你需要手写一个高并发socket服务而不引入框架可以用selectors实现import socket import selectors sel selectors.DefaultSelector() def accept(sock): conn, addr sock.accept() conn.setblocking(False) sel.register(conn, selectors.EVENT_READ, read) def read(conn): data conn.recv(1024) if data: print(data) else: sel.unregister(conn) conn.close() server socket.socket() server.bind((127.0.0.1, 9999)) server.listen(100) server.setblocking(False) sel.register(server, selectors.EVENT_READ, accept) while True: events sel.select() for key, mask in events: callback key.data callback(key.fileobj)这段代码的核心思路是事件循环本身是唯一的循环每个连接注册一个回调当内核告诉你“连接上来了”或“数据到了”才去操作几乎没有空闲等待浪费。5.3 如何根据业务选择正确的 IO 模型我会按业务类型做选择脚本要批量请求几十个API直接用asyncioaiohttp轻量高效。需要处理大量长连接聊天、推送、网关用asyncio它能挂住几十万连接不炸。大量小文件读写线程池 队列因为文件IO的不可控因素多协程收益有限。大文件高吞吐传输用多进程或独立的异步IO库必要时上io_uring。IO模型的选型不是孤立的。只要网络IO密集协程几乎永远是第一选择只要涉及CPU计算就得靠多进程混合场景则要组合拳。记住一条原则让CPU尽量别空等让等待尽量不占线程。6. 实战排障记录我踩过的并发坑6.1 Windows 下进程被锁定的排查思路热搜词里有“f盘被另一个进程锁定”“未找到baidunetdiskhost进程”“msedgewebview2.exe进程如何关闭”这类题。这类问题的本质是Windows文件句柄被某个进程占用导致其他进程无法访问文件。碰到这种情况我的排查路径是用微软官方工具handle.exe查看哪个进程占用了指定路径的句柄。或者打开资源监视器在“CPU”页的“关联的句柄”搜索文件路径马上就能看到占用的进程。确认占用进程后判断该进程是否还在干活。不是所有占用的进程都能直接关比如系统关键进程瞎关会蓝屏。这个问题的启示是写多进程程序要认真管理文件句柄。子进程继承父进程已打开的文件句柄可能导致父进程想释放文件却释放不掉因为子进程还攥着不放。所以创建子进程前尽量先关掉不需要继承的文件。6.2 线程池配置不当引发的 CPU 100%我曾维护过一个爬虫脚本ThreadPoolExecutor设了32个线程结果部署后CPU直接飙到100%。直觉以为是目标网站响应慢但抓线程栈才发现问题。用py-spy dump --pid pid看到所有线程都卡在同一个位置一个while True循环里反复执行queue.get_nowait()。当队列为空时get_nowait()立刻抛异常线程马上再循环等于32个线程在疯狂空转。解决方案是换成阻塞式queue.get(timeout0.5)让线程在队列为空时真正睡一会儿CPU立刻降下来。代码层面的教训是并发任务里尽量不要用忙轮询的方式等待资源。queue.get()默认就会阻塞等待没有必要自己写循环。如果担心阻塞无法退出就设置合理的超时时间。6.3 进程堆大小与内存问题的关系热搜词里那条“进程堆大小调整为8000还是报错java: java.lang.outofmemoryerror”虽然说的是Java但原理同样适用于Python。很多人以为把内存调大就能解决OOM实际上Python里更常见的是多进程场景下每个子进程各自加载一份完整的大对象内存随进程数线性增长。举个例子一个进程池有8个worker每个worker加载一个1GB的词典数据那总内存就是8GB机器再大也扛不住。优化思路有几种改用父进程加载数据后通过共享内存传给子进程。把大词典落到MMAP文件里子进程按需读取。减少进程数配合协程或异步IO分担IO密集型操作。另一个容易忽视的地方是multiprocessing传递大对象给子进程时spawn模式会重新序列化并复制一份数据内存瞬间翻倍。如果数据实在太大宁可反复从磁盘读取也不要一次性塞给每个子进程。6.4 守护进程与会话热词里有“守护进程与会话”。很多同学觉得daemonTrue就是把程序变成系统守护进程其实不对。Python里的Process.daemonTrue只是说这个进程随父进程退出不会让程序变成真正的后台服务。要在生产环境跑后台Python服务用systemd或supervisor管理才是正路它会处理进程组、会话、日志捕获和自动重启。排查“未找到xxx进程”时第一件事不是怀疑进程被杀了而是看它的启动方式。用ps aux看看进程的TTY列如果是?说明它可能是被某个服务管理器拉起的如果找不到再看dmesg | grep -i oom很可能内存不够被OOM killer干掉了。进程是否存活只是表象背后的崩溃原因才是关键。7. 并发选型决策表与心得7.1 一张表搞定并发方案选型很多情况下选型定下来问题就解决了一半。我根据自己的实践经验做了这张决策表可以作为参考。场景推荐方案注意事项CPU密集计算ProcessPoolExecutor进程数约等于CPU核心数避免子进程重复加载全局数据网络IO密集大量请求asyncio 异步HTTP客户端用Semaphore控制并发量不要用同步requests库文件/数据库IO密集ThreadPoolExecutor线程数压测确定不要盲信公式高并发长连接服务asyncio selectors注意事件循环里不能有同步阻塞操作大量短期小任务ThreadPoolExecutor复用线程池避免反复创建线程需要强隔离的任务多进程IPC成本高先评估数据交换量任务间共享大量状态threading Lock锁粒度控制很重要考虑用队列代替共享变量这张表背后还有一条隐含原则没有银弹。进程、线程、协程不是互斥的真正复杂的系统往往是三者混用。关键是让每种机制都在它最适合的层次上工作。7.2 个人心得并发编程的核心不是“并行”而是“协调”写并发代码这几年我最大的体会是并发编程的难点从来不是“创建多少个执行单元”而是“怎么协调它们之间的节奏”。就像城市交通不是车越多越快而是信号灯、车道规划、公交车专用道互相配合。Python给我们的选择足够多进程负责隔离线程负责共享与并行等待协程负责超高并发下的轻量切换IO多路复用负责解决“如何高效等待外部事件”。最后再分享一个对我帮助很大的习惯每次写并发代码先画一张任务流标清楚哪些步骤可以并发哪些步骤必须串行哪些数据需要共享哪些只能隔离。这张图不需要画得多规范但画完后你会发现自己对选型的判断准确很多。并发本身不产生价值正确的协调才产生价值。
返回列表