Python 并发编程实战:线程、进程与协程的选择策略

发布时间:2026/8/14 16:17:49
Python 并发编程实战:线程、进程与协程的选择策略 关于并发编程存在三条可选途径, 其一为多线程形式, 其二属多进程方式, 其三是异步协程模式。表面上看这三种方式似乎均可实现同时进行多项事务的效果, 然而它们各自所适用的具体场景的差异程度极大。一旦选错了相应的模型, 在代码已经全部编写完成之后才察觉到性能不但没有提升反而出现下降的情况, 像这类案例在代码领域当中实在是太过常见了。这段文字起始于GIL, 进而逐个剖析三种并发模型的底层机制, 接着阐述其适用边界, 随后展示代码示例, 最终给出一张清晰明确的选型决策表。一、GIL绕不开的话题一个解释器的设计决策是GILLock即全局解释器锁, 在同一时刻, 仅有一个线程可执行字节码那这究竟意味着啥呢? 哪怕你开启了有八个线程, 让它们运行在具备八核的中央处理器之上, 然而代码的实行依旧是呈串行化进行的。多个线程在处于纯粹由中央处理器操作带动负担集中执行的任务之时, 并不能在速度上演进的节奏获得提升, 反而会更缓慢, 究其原因在于线程更迭转换带来的额外消耗所导致的。import threading import time def cpu_bound(n): total 0 for i in range(n): total i ** 2 return total start time.perf_counter cpu_bound(10_000_000) print(f单线程: {time.perf_counter - start:.3f}s) start time.perf_counter threads [threading.Thread(targetcpu_bound, args(2_500_000,)) for _ in range(4)] for t in threads: t.start for t in threads: t.join print(f4线程: {time.perf_counter - start:.3f}s) # 大概率比单线程更慢不过, GIL存在着一个关键的例外情况, 那就是, 当线程处于等待IO这种状态的时候, 这里所说的IO涵盖网络请求、磁盘读写以及数据库查询, 在这个时候, GIL是会被释放掉的, 然后其他的线程就能够继续去执行了。而这恰恰就是多线程在那种IO密集型场景之下依旧能够生效的最为根本的原因所在了。二、IO 密集场景的首选 2.1 基本用法import threading import requests from concurrent.futures import ThreadPoolExecutor urls [https://httpbin.org/delay/1] * 10 def fetch(url): resp requests.get(url, timeout10) return resp.status_code with ThreadPoolExecutor(max_workers5) as executor: results list(executor.map(fetch, urls)) print(f完成 {len(results)} 个请求)2.2 线程安全与锁多线程共享内存空间写操作必须加锁counter 0 lock threading.Lock def increment: global counter for _ in range(100_000): with lock: # 等同于 lock.acquire / lock.release counter 1看起来 1 好像是原子操作, 然而在字节码层面它却被分成了多条指令, 先是读取, 然后加 1, 最后写入。当不加锁的时候, 线程 A 会读到旧的值, 线程 B 同样也会读到旧的值, 之后它们各自进行加 1 操作并写回, 最终导致结果少算了一次。和queue.Queue不同的另一种选择是 , 它是线程安全的队列 , 这种队列出于自然的特性而适合生产者 - 消费者模式:from queue import Queue import threading def producer(q): for i in range(100): q.put(i) def consumer(q): while True: item q.get if item is None: break process(item)2.3 适用场景 三、CPU 密集场景的解法把进程进行更换从而绕开了GIL, 每一个进程具备着独立的, 作为独立存在的解释器以及占据各自固定范围的内存空间, 切实地凭借多核CPU去达成相关任务。from concurrent.futures import ProcessPoolExecutor import time def cpu_bound(n): total 0 for i in range(n): total i ** 2 return total start time.perf_counter with ProcessPoolExecutor(max_workers4) as executor: results list(executor.map(cpu_bound, [2_500_000] * 4)) print(f4进程: {time.perf_counter - start:.3f}s) # 接近线性加速代价也很明显3.1 进程间通信from multiprocessing import Process, Queue def worker(q): q.put({result: 42}) if __name__ __main__: q Queue p Process(targetworker, args(q,)) p.start result q.get # {result: 42} p.join对于大量数据的共享3.8 给出了更为高效的途径, 规避序列化所带来的开销。四、高并发 IO 的现代方案的核心思想为协作式多任务, 在单个线程内部, 借助await去实现显式地交出控制权, 从而能够使得事件循环于同一线程中对多个协程予以调度。import asyncio import aiohttp import time async def fetch(session, url): async with session.get(url) as resp: return await resp.text async def main: urls [https://httpbin.org/delay/1] * 50 async with aiohttp.ClientSession as session: tasks [fetch(session, url) for url in urls] results await asyncio.gather(*tasks) return results start time.perf_counter asyncio.run(main) print(f50个请求耗时: {time.perf_counter - start:.3f}s) # 约 1-2 秒50个请求, 每个请求延迟1秒, 按同步顺序执行的话, 需要50秒, 而并发执行的话, 大概只需1到2秒, 这取决于带宽以及连接池大小。这便是其核心优势: 以极低的线程开销处理海量的IO并发。4.1 async/await 的使用边界async def fetch_data: data await db.query(SELECT ...) return data asyncio.run(fetch_data) def main: data await fetch_data # SyntaxError4.2 同步代码与异步代码的桥接于实际项目里, 并非能够使所有代码皆为async。于异步的上下文环境之中, 去对同步函数展开调用import asyncio from concurrent.futures import ThreadPoolExecutor executor ThreadPoolExecutor(max_workers4) async def main: loop asyncio.get_running_loop result await loop.run_in_executor(executor, blocking_io_function)反转过来, 于同步代码里调用异步函数之时, 要使用 .run , 然而需留意, 不要进行嵌套, 也就是 .run 万万不可在那样一个业已有事件循环的线程当中再次去调用。五、并发模型选型决策场景 推荐方案 原因IO 密集网络请求、文件读写、数据库查询GIL 在 IO 等待时释放CPU 密集数值计算、图像处理、加密绕过 GIL利用多核高并发 IO万级连接如 、长轮询单线程事件循环无上下文切换开销中小规模 IO 并发百级已有同步代码库改造成本最低混合场景IO CPU 均有异步调度 CPU 密集型子任务需要共享大量内存状态进程间共享内存成本太高六、一个混合场景的实战示例设想存在这样一个任务, 从多个API那里拉取数据呢, 这属于IO密集型的, 之后针对每一条数据, 去做CPU密集型的解析。import asyncio from concurrent.futures import ProcessPoolExecutor def heavy_parse(raw_text: str) - dict: CPU 密集的解析逻辑。 return parsed async def fetch_and_parse(session, url, executor): async with session.get(url) as resp: raw await resp.text loop asyncio.get_running_loop return await loop.run_in_executor(executor, heavy_parse, raw) async def main(urls): executor ProcessPoolExecutor(max_workers4) async with aiohttp.ClientSession as session: tasks [fetch_and_parse(session, url, executor) for url in urls] results await asyncio.gather(*tasks) executor.shutdown return results存在一种模式, 它综合了具备高并发IO能力的某些特性, 以及拥有 CPU并行能力的特定方面, 这种模式适用于数据管道场景, 适用于ETL场景, 适用于批量API处理等场景。最后关键在于明确任务性质来进行并发选型, 任务的瓶颈究竟是在IO方面, 还是在CPU方面并发规模到底是百级别的, 还是万级别的? 现有代码到底是同步的, 还是异步的?几条原则

相关新闻