ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Python多线程并发编程实战:GIL、线程池与锁的避坑指南

Python多线程并发编程实战:GIL、线程池与锁的避坑指南 做 Python 后端或者爬虫的几乎都会碰到“多线程”。说实话Python 的多线程被议论得最多有人说它没用因为 GIL 锁卡死 CPU 密集型任务也有人说它在 IO 密集型场景里能提速好几倍。谁对谁错我在实际项目里测过、踩过、也把生产环境的线程池调优过一遍。这篇东西不是来复述文档的而是把我这几年用 Python 多线程处理任务的经验拆开讲清楚包括什么时候用它、怎么设计任务队列、锁到底怎么加才不踩坑以及我碰过的各种诡异问题。Python 多线程最容易被新手误用的点是把“多线程”当成“多核并行”。在 CPython 解释器里GIL全局解释器锁决定了同一时刻只有一个线程在执行 Python 字节码所以如果你想靠线程把四核 CPU 全部跑满基本没戏。但多线程在 IO 密集型场景下又是绝对的王者比如爬虫、下载文件、读写数据库、调用接口因为线程在等待 IO 返回时会让出 GIL这时候你开的几十个线程就能把阻塞时间重叠起来办事效率完全不一样。这篇文章适合这么几类人看写爬虫想提速的、做接口服务想降低响应时间的、刚把 Python 当主语言想搞清楚并发模型的以及在多线程下有数据错乱或死锁困扰的。我不会只给结论每个判断都会给场景和实测经验看完了可以直接抄到你的项目里。1. 为什么 Python 多线程总被误解以及它真正适用的两类场景1.1 GIL 并没有让多线程一无是处先把这个最核心的原理讲透。GIL 的全称是 Global Interpreter Lock它是 CPython 内存管理机制的一部分。因为有它Python 对象的引用计数才不用额外做大量的线程锁保护内存安全实现起来简单得多。代价也很明确就是同一个进程里真正的计算只能一个线程来跑。但注意“同时只有一个线程执行 Python 字节码”这句话有一个隐藏前提这段时间是“执行 Python 代码”的时间。当线程在做阻塞式 IO比如requests.get()等待响应、time.sleep()、socket.recv()它并没有在运行 Python 字节码锁会在这时释放让别的线程去跑计算。所以多线程在 IO 密集场景下看起来就像是“并行”的因为这个过程中线程的可执行时间被调度器充分塞满了。理解了这个原理你就能解开一个新手最常有的困惑为什么明明包里写了threading.Thread也用了start()程序执行时间几乎没变化因为任务要么是纯 CPU 计算要么线程数量太少IO 重叠得不够充分。1.2 快速判断你的任务适不适合上多线程我给团队做代码审查时判断一个任务适不适合多线程就三板斧任务是否在等待外部资源例如 HTTP 请求、数据库查询结果、磁盘文件读取、外部进程输出。是的话适合。任务是否主要是纯计算例如大规模的数值循环、图像像素逐点处理、字符串规律匹配。是的话优先考虑multiprocessing或重写相关库去用 C 扩展。任务是否有大量数据共享和跨线程通信如果共享数据极其复杂写锁比写线程还难也许该用无锁队列或消息传递。现实里很多任务都是混合型的。比如爬虫既要下载网页IO又要解析文本CPU 轻度可能还要写库IO。这种情况下用多线程完全没问题因为整体耗时主要被 IO 支配计算占比小GIL 基本不会成为瓶颈。实测一个例子用单线程爬 500 个页面假设每个请求耗时 0.2 秒总耗时是 100 秒。开 20 个线程后理想情况下是 5 秒左右算上排队和网络波动实际也就在 7~10 秒之间。这个量级的提升是非常直观的。2. 线程的三种创建方式该用哪个心里要有数2.1 手工创建threading.Thread小而直观最基础的方式是直接继承或直接实例化Thread。比如import threading import time def worker(name): time.sleep(1) print(f{name} 完成) threads [] for i in range(5): t threading.Thread(targetworker, args(f任务-{i},)) threads.append(t) t.start() for t in threads: t.join()这段代码没什么坏处但有个明显体验问题你需要自己维护线程列表、自己join()等所有任务结束。一旦任务数量变成几百几千这种手工管理就显得落后了。它还缺少一个很重要的机制——任务队列和结果的回收。你用Thread启动一个线程想知道它返回结果还得自己定义全局变量或者把结果塞进Queue非常啰嗦。所以我的看法是threading.Thread适合写脚本、临时验证想法或者实现长驻后台线程时用比如一个定时上报系统状态的守护线程。正式处理批量任务时我会直接跳到线程池。2.2ThreadPoolExecutor现代项目首选写 Python 多线程处理任务我大部分情况用的是concurrent.futures.ThreadPoolExecutor。它帮我们管理线程生命周期、限制最大并发数、回收结果、处理异常代码也能精简不少。from concurrent.futures import ThreadPoolExecutor, as_completed import requests urls [fhttps://example.com/page/{i} for i in range(100)] def fetch(url): r requests.get(url, timeout10) return r.status_code, len(r.content) with ThreadPoolExecutor(max_workers20) as pool: futures [pool.submit(fetch, url) for url in urls] for future in as_completed(futures): status, size future.result() print(status, size)max_workers不是越大越好。开线程有内存开销切换线程也有成本。实测中 IO 密集任务线程数设成目标 IO 并发数的 1.5~3 倍比较合适如果任务里有大量阻塞等待比如接口平均耗时几秒那可以开多一点但一般也不建议超过 100除非你真的清楚系统资源上限。as_completed是处理结果的利器。哪个任务先完成就先取哪个的结果。这样处理日志、入库、重试都特别自然不用等最慢的那个拖后腿。2.3Pool.map和start_join一些另类但好用的用法除了submitThreadPoolExecutor还支持map。它和内置map长得很像只是并发执行with ThreadPoolExecutor(max_workers8) as pool: results pool.map(fetch, urls)map返回的结果顺序和传入顺序一致所以如果你希望“结果按提交顺序拿到”这个就比as_completed方便得多。代价是如果前面的请求很慢后面即使已经完成了也要等前面的结果先返回。还有一个小众但实用的场景把多个线程的启动和结束封装成上下文管理器。ThreadPoolExecutor本身支持上下文管理协议with块结束时会自动调用shutdown(waitTrue)保证所有线程都跑完才退出。这一点非常稳你不需要自己写join()循环。3. 线程安全、共享变量与锁的正确姿势3.1 千万别直接在线程里改共享列表新手最常见的翻车现场是这种代码import threading results [] def work(i): results.append(i * 2) threads [threading.Thread(targetwork, args(i,)) for i in range(10)] for t in threads: t.start() for t in threads: t.join() print(len(results)) # 可能不是 10为什么结果数量不定因为results.append()操作不是原子的。它要做两步找到列表的尾部内存位置写入新元素并修改长度。两个线程同时执行时后写的数据可能覆盖前一个线程的索引导致结果丢失或者长度不对。GIL 只能保证单个字节码指令的原子性不能保证 Python 对象操作整体安全。实际生产环境里这种问题早期几乎不会暴露。本地跑一遍可能结果刚好正确但稍微增加并发量、换个机器性能数据就丢了。排查起来还特别困难因为不是每次都发生。解决办法很简单用threading.Lock保护临界区或者直接用queue.Queue。3.2 用queue.Queue替代手工锁Python 的queue.Queue内部已经做好了线程同步。它使用一个由threading.Lock和Condition组合起来的机制保证put()和get()是线程安全的。所以当你想让多个线程协作完成任务时与其共享一个列表再加锁不如把任务塞进队列里。import queue import threading import random task_queue queue.Queue() def producer(): for i in range(10): task_queue.put(i) print(f生产任务 {i}) def worker(): while True: try: task task_queue.get(timeout3) except queue.Empty: break print(f处理任务 {task}) task_queue.task_done() prod threading.Thread(targetproducer) prod.start() workers [threading.Thread(targetworker) for _ in range(3)] for w in workers: w.start() prod.join() workers[j].join() task_queue.join()task_done()和join()的配合是理解这个代码的关键。消费者每从队列里取走一个任务处理完之后必须调用task_done()表示这个任务完成了。主线程调用task_queue.join()时会阻塞直到队列里的所有任务都调用了task_done()这样你就不用自己统计完成次数了。这种“生产者-消费者”模型是大多数多线程任务系统的基础。爬虫里生产者负责生成待抓取的 URL消费者负责下载和解析中间的队列还能做顺便做去重和限流。3.3 锁的类型、超时和死锁防护并不是所有场景都能用队列替代锁。很多时候你确实需要共享一个状态变量比如统计失败次数、记录全局配置变化。这时就要用锁。threading.Lock()是最基础的互斥锁。同一时刻只能有一个线程持锁其他线程会被阻塞。threading.RLock()是可重入锁允许同一个线程多次acquire()而不死锁。递归函数里如果用了Lock容易把自己锁死所以这种情况要用RLock。锁的一个大坑是死锁。典型场景是线程 A 持有锁 1想再拿锁 2而线程 B 持有锁 2想再拿锁 1。两个线程互相等待谁也别想跑。规避方法有很多尽量缩小加锁范围只在修改共享数据时加锁不要在锁里做 IO 耗时操作。制定锁的获取顺序。比如所有地方都先拿锁 1 再拿锁 2就可以避免循环等待。使用with lock:语句而不是手动acquire()/release()。with块即使异常退出也会释放锁减小了忘释放的概率。如果确实需要在等待锁时有超时底线用lock.acquire(timeout3)拿不到锁就走另外的业务逻辑不硬等。我还会强烈建议不要把持锁时间拉长去调用外部 API。你在持锁时其他所有需要这把锁的线程都会卡住等于把多线程变回了单线程而且如果外部接口响应慢线程还会越积越多最终拖垮整个进程。4. 线程池与任务队列结合的真实项目设计4.1 如何设计一个可控的多线程任务消费系统我看过很多代码一上来就开 200 个线程没人管资源、没人管重试、没人管任务堆积。等到服务突然卡死查了半天才发现是任务队列里积压了几万条数据内存直接爆掉。设计一个可控的多线程任务系统至少有这五件事要到位并发上限可配置。不要写死用环境变量或配置中心下发。队列长度有界。可以用queue.Queue(maxsize500)满了之后put()会阻塞这样下游处理不完时上游自然放慢而不是无限堆积。任务要有超时。线程正常执行时很难取消但你可以让任务内部设置超时或者用future.result(timeout10)来主动放弃一个过慢的结果。失败任务要重试但要有重试上限。很多人忘了这一点导致失败任务反复执行把数据库都打爆。有退出机制。进程收到终止信号时要能优雅地停止接收新任务等待已执行的任务完成再退出。4.2 实战用 ThreadPoolExecutor 写一个可复用的调度模块下面这个模块是我在项目里用了很久的一个简化版本。它可以投递任务、回调结果、控制最大并发、收集异常。import logging import time from concurrent.futures import ThreadPoolExecutor, as_completed from queue import Queue, Full from typing import Callable, Iterable, Any logger logging.getLogger(__name__) class TaskDispatcher: def __init__(self, max_workers10, queue_size200): self.pool ThreadPoolExecutor(max_workersmax_workers) self.queue Queue(maxsizequeue_size) self.running False def submit(self, fn: Callable, *args, **kwargs): if self.queue.full(): # 队列已满丢弃或阻塞按需调整 logger.warning(task queue is full, drop the task) return future self.pool.submit(fn, *args, **kwargs) self.queue.put(future) future.add_done_callback(self._on_done) return future def _on_done(self, future): try: result future.result() logger.info(task success: %s, result) except Exception as exc: logger.error(task failed: %s, exc) finally: self.queue.task_done() def wait(self): self.queue.join() def shutdown(self): self.pool.shutdown(waitTrue)用的时候只需要dispatcher TaskDispatcher(max_workers20, queue_size500) def load_data(item_id): time.sleep(1) return {id: item_id, status: ok} for i in range(1000): dispatcher.submit(load_data, i) dispatcher.wait() dispatcher.shutdown()这里有个细节add_done_callback是在任务完成的线程里执行的不是主线程。所以在回调里更新 UI 或者写日志时要小心如果回调里又访问了同一个可变对象仍然要保证线程安全。4.3 并发数怎么定给一个可参考的经验公式并发参数最让我头疼因为网上各种说法都有。实测下来IO 密集场景可以这么估算线程数 ≈ 单任务耗时内的等待时间比例 × 目标吞吐量 × 一个冗余系数。举个例子一个接口平均响应 0.1 秒我希望每秒完成 50 个请求那么平均每个任务占用连接的时间是 0.1 秒需要同时处理的连接数是 50 × 0.1 5所以线程数 5 就够了。但考虑到网络不稳定、某些请求会变慢我会乘以 2 到 3最终开 10~20 个线程。不要迷信“开得越多越快”。线程本质上是操作系统资源每个线程都有栈空间默认大概 8MB 虚拟内存和调度成本。100 个线程的纯 Python 任务GIL 竞争就会非常明显1000 个线程系统光切换上下文就忙不过来速度不升反降。如果任务里有非常耗时的阻塞点比如下载一个 1GB 文件那并发数可以适度调大一些如果任务本身几毫秒就完成再开高并发意义不大因为创建线程和调度的开销都快要超过任务本身了。5. 多线程与多进程怎么搭什么时候必须换方案5.1 CPU 密集任务为什么用多进程更香GIL 对 CPU 密集任务的限制是硬伤。比如做图片批量滤镜、解析大量 JSON 并写回这类任务不涉及等待外部资源纯靠 CPU 计算。你用多线程跑起来会发现 CPU 总使用率一直停留在 100% 左右因为你只有一个 Python 线程在真正计算其他线程在等锁。这时候正解是concurrent.futures.ProcessPoolExecutor或多进程。每个进程都有独立的 Python 解释器和独立的 GIL理论上可以跑满多核。用起来和ThreadPoolExecutor几乎一样from concurrent.futures import ProcessPoolExecutor def cpu_intensive(n): return sum(i * i for i in range(n)) with ProcessPoolExecutor(max_workers4) as pool: results list(pool.map(cpu_intensive, [10000000] * 8))但要注意多进程之间不能直接共享普通 Python 对象数据需要通过序列化传输通信成本比线程高不少。所以如果你的任务是“IO 为主、计算很少”用多进程反而没必要序列化和进程切换的消耗会让收益变负数。5.2 一个任务混合了 IO 和 CPU怎么拆真实场景里很多任务是混合的。比如爬虫下载完页面后要做 HTML 解析解析是 CPU 计算下载是 IO。这种任务如果只开线程解析阶段会被 GIL 卡住如果只开进程下载阶段的并发连接数又不好控制。我的方案是分层主进程用ThreadPoolExecutor负责任务调度和并发 IO真正吃 CPU 的解析操作单独丢给一个ProcessPoolExecutor。可以用“线程池取任务 - 子进程池计算 - 主线程收集结果”的流水线模型。网上也把这叫做“混合并发模型”。简单示意from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor def download(url): return url, bhtml... def parse(html): return len(html) with ThreadPoolExecutor(max_workers20) as tpool, ProcessPoolExecutor(max_workers4) as ppool: futures [] for url in urls: f tpool.submit(download, url) futures.append(f) for f in futures: url, html f.result() p ppool.submit(parse, html) print(p.result())这个模型保证了下载阶段的高并发不被计算任务卡住也保证了计算阶段能用到多核。代价是代码复杂度上升而且ProcessPoolExecutor在 Windows 上有一些注意点比如必须在if __name__ __main__里创建进程池否则会递归创建进程导致报错。5.3 什么时候该考虑 asyncio既然多线程在 IO 场景不错为什么还要有 asyncio因为在超高并发连接场景线程本身的资源开销可能是最大的瓶颈。比如一个服务要同时保持几万个 WebSocket 或 HTTP 长连接开几万个线程显然不现实而 asyncio 是单线程事件循环内存占用小得多可以轻松支撑几十万连接。如果你的任务是以 IO 等待为主、并且你能把代码改成async/await风格那 asyncio 是更优解。但它的学习曲线比较陡而且和同步库混用会有额外麻烦。很多新手问我“多线程爬虫是不是最好的方案”我一般会说中小规模爬虫用多线程开发快、调试容易真要大规模高并发采集再考虑 asyncio 或者上面说的混合方案。没有银弹看场景选。6. 我踩过的坑多线程调试、退出与资源泄漏6.1 异常不是你看不见就不会炸多线程里有一个特别折磨人的特点线程里抛出的异常不会像单线程那样直接把整个程序打断而是被线程吞噬掉默认打印一段 traceback程序继续跑。这导致一个问题业务逻辑里某一步经常失败但你没有感知直到统计报表数据对不上才发现。正确做法是在线程入口函数的最外层包一个 try/except把异常记录到日志或者队列里。使用ThreadPoolExecutor时future.result()会把异常再抛出来所以你取结果的地方也要有异常处理。不要图省事忽略异常调试成本会翻好几倍。6.2 主线程退出子线程还在干活怎么办默认情况下线程不是守护线程主线程结束后进程会等待所有非守护线程结束。如果你用了ThreadPoolExecutorwith块会等待所有任务完成通常没问题。但如果你手写Thread并且子线程里有一个无限循环主线程想退出时会发现卡住退不出去。解决办法是设置daemonTrue或者自己实现一个停止标志事件。import threading import time stop_event threading.Event() def worker(): while not stop_event.is_set(): time.sleep(1) print(working...) t threading.Thread(targetworker, daemonTrue) t.start() time.sleep(3) stop_event.set()用Event比用一个全局布尔变量更严谨。因为Event本身是线程安全的而且在多个线程等待通知时表现更好。如果你的线程里有queue.get()阻塞等待想停止时要么用timeout定期醒来检查事件要么在 queue 里放一个特殊消息让它自行退出。6.3 资源泄漏连接池、文件句柄和线程堆栈多线程程序还有一个看不见的坑就是资源泄漏。每个线程打开一个数据库连接、一个文件句柄任务跑完又不关闭时间一长连接池耗尽、文件句柄数超过系统上限程序就会莫名其妙报OSError: [Errno 24] Too many open files。排查这类问题我建议在开发环境做一次压力测试把并发数调到线上配置的 2 倍跑一段时间用lsof -p pid | wc -l看文件句柄数是否持续增长。如果只增不降基本就能确定存在资源泄漏。解决的核心思路是把“每次任务创建连接”改成“使用连接池”并且一定用with或try/finally来保证释放。数据库连接池在SQLAlchemy、redis-py、requests.Session这些库里都有现成的支持没必要自己造。requests.Session本身就是线程安全的但官方文档也说了它不保证绝对的线程安全推荐每个线程用独立的 Session或者用requests配合urllib3的连接池机制。6.4 调试多线程的心得加 trace_id 和顺序日志多线程程序难调试主要是因为日志交错、上下文丢失。我现在的习惯是每个任务开始和结束时都打一条日志带上任务 ID 和上下文信息比如import logging import threading from uuid import uuid4 def run_task(task_id): thread_id threading.get_ident() logging.info(task %s started on thread %s, task_id, thread_id) try: # do something ... except Exception: logging.exception(task %s failed on thread %s, task_id, thread_id) else: logging.info(task %s finished on thread %s, task_id, thread_id)配合 Traceback 和上下文 ID你可以在日志系统里把同一个任务的全部日志拉出来定位问题就快得多。这也算是我踩坑之后形成的固定习惯。7. 常见问题排查速查表这里整理一下我在论坛和团队里被反复问到的多线程问题每个都附上排查思路和解法。现象可能原因排查与解法任务全部执行完但join()一直不返回线程里有阻塞等待代码没退出或者在循环中持续新建线程导致僵尸线程堆积检查线程函数是否有死循环加入停止条件确认线程数量是否在增长多个线程修改同一个变量后数据不对缺少锁保护多个线程同时写共享对象使用Lock或改用queue.Queue避免共享可变对象多线程程序反而更慢任务本身是 CPU 密集GIL 导致线程切换开销换ProcessPoolExecutor或使用numpy等释放 GIL 的扩展库某个线程抛了异常主线程完全不知道异常被线程吞噬没有向外传递在future.result()里捕获异常或在线程函数内手动上报程序运行一段时间后文件句柄耗尽每个线程创建了连接或文件没有正确关闭使用连接池确保with或finally释放资源ThreadPoolExecutor卡住不退出线程池中还有任务在跑或者任务内部有死锁检查任务是否有无限等待可以加shutdown(waitFalse)并打印堆栈进程卡死无法 CtrlC 退出非守护线程还在运行且KeyboardInterrupt只在主线程生效子线程设置daemonTrue或实现优雅退出逻辑这个表看着简单但绝大多数多线程问题逃不出这几类。解决问题最快的路径不是看报错而是先明确两个问题现在有多少线程在跑它们在等什么资源把这俩搞清楚了问题基本就定位了一半。8. 最后的几个实操建议我自己的经验是多线程方案的成败主要看任务拆分得清不清楚而不是线程写得多花哨。一个很简单的任务比如批量下载文件如果任务列表完整、每个任务之间没有依赖关系那用ThreadPoolExecutor只需要几行就能写得很稳。相反如果任务之间有依赖关系后一个任务要用前一个任务的结果那就要认真设计任务 DAG 或者分层队列了。这个时候千万别图省事把所有逻辑塞进一个线程函数里到时候调试会非常酸爽。另外一个建议是上线前一定要加一个“限流开关”。多线程提速是好事但如果不控制整体请求速率很容易把下游接口打挂然后触发封禁、限流或者生产事故。我习惯在任务代码里加上一个信号量或令牌桶比如每秒最多放行多少个任务去执行这样既保证了吞吐量又不会对外部系统造成冲击。threading.Semaphore就能做简单的限流或者用rate-limiter这类库。最后想啰嗦一句如果你真的把多线程调优弄得很深入就会发现大多数时候瓶颈不在“线程”而在“等待”。优化多线程程序先看等待时间花在哪里再去想线程数量、锁粒度、任务分配方式。只要抓住这个核心线程池方案很难出大错。我至今还记得第一次用线程池把爬虫从一小时压到三分钟时的那种成就感。多线程这个东西说难也难说简单也简单无非就是搞清楚并发模型、线程安全和资源管理这三件事。希望这篇内容能让你少走一些弯路也欢迎你把自己碰到的奇葩多线程问题丢到评论区一起聊一聊。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进