
Python线程同步这个话题网上的教程分成两个极端要么只讲threading.Lock怎么用配一个最简单的计数器例子就草草收场要么一上来就搬出GIL、GIL、GIL最后得出结论“反正有全局锁多线程就是个摆设”。这两个极端我都经历过都会把人带偏。我写这篇文章是想把线程同步这件事真正讲透——锁为什么存在、锁该怎么选、线程池和同步机制怎么搭配、死锁到底怎么定位和避免。内容基于我在实际项目中踩过的坑目标读者是写过一些多线程代码但总感觉“不太稳”的Python开发者。看完你至少能回答三个问题这段共享数据要不要加锁、该用哪种同步原语、线程池的任务结果怎么安全拿回来。另外先划个界限很多人把线程同步和数据库同步、硬件时钟同步、文件同步软件混为一谈。那些是另一个领域的问题本文只聊Python多线程编程里的同步机制。下面所有内容都基于CPython 3.8版本其他Python解释器比如PyPy行为略有差异但同步思路完全通用。1. 为什么需要线程同步1.1 竞态条件是怎么发生的先看一个最经典的例子两个线程各自对同一个变量执行1000次自增。直觉上结果应该是2000但实际跑出来经常是1793、1861这种数字每次还不一样。原因在于count 1这行代码在Python字节码层面根本不是一步操作。你可以自己用dis模块看一眼import dis def add_one(count): count 1 return count dis.dis(add_one)输出会显示count 1被拆成了LOAD_FAST、LOAD_CONST、INPLACE_ADD、STORE_FAST这四步。两个线程可以同时读到一个旧值各自加完再写回去结果就丢了一次更新。这就像两个人同时在一本账本上登记支出A看到余额100记了100消费变成0B也在同一时刻看到100记了100消费也变成0账本就错了。这种多个线程因为交错执行导致的错误结果就是竞态条件。线程同步要解决的本质上就是让这种不可控的线程调度变得可控。你不能依赖运气必须用机制保证“读-改-写”这个复合操作不被别的线程打断或者至少让共享状态的变化可以被其他线程正确感知。1.2 GIL是不是免死金牌不少人认为CPython有GIL全局解释器锁同一时刻只有一个线程在执行Python字节码那共享变量是不是就安全了这是对GIL最大的误解。GIL确实保证了“字节码指令”级别的原子性但它不保证“业务操作”级别的原子性。上面那个count 1被拆成四步执行线程可能在任意两步之间被切换出去。GIL只保证每一步本身不被打断但不管你的多步操作是否连续完成。更关键的是GIL会在两种情况下主动释放线程等待I/O操作时以及执行某些C扩展代码时。所以当你写网络请求、文件读写、数据库查询时GIL是放开的其他线程有机会执行。你辛苦维护的数据结构如果没有加锁分分钟被钻空子。举一个真实业务场景一个电商系统的库存扣减。两个用户同时下单各自线程都读取到库存为1都判断“还有货”都执行扣减最终库存变成-1。GIL救不了这种问题必须依靠同步机制来保证“读库存-判断-扣减”整个流程的互斥。1.3 同步的本质让不可控的调度变得可控理解了竞态条件和GIL的局限性你就能明白同步机制设计的出发点。锁、信号量、事件、条件变量、队列这些工具统统是干嘛的一句话协调多个线程对共享资源的访问顺序和时机。锁负责互斥同一时刻只能有一个线程进入临界区信号量负责限流同时最多N个线程访问资源事件负责通知一个线程等另一个线程的信号条件变量负责更精细的等待与唤醒队列则把共享数据藏在一个内部有锁的容器里让你写出天然的线程安全代码。我个人在实际项目里的一个粗浅经验是能用队列传递数据就不要自己手写共享变量加锁必须共享状态时优先用锁保护最小范围的临界区需要多个线程协作处理复杂状态变化时才考虑条件变量和事件。这个经验帮我避免了很多线上事故。2. 六大同步原语逐个拆解2.1 Lock与RLock最基本的互斥控制threading.Lock是最基础的互斥锁。它有两种状态锁定和未锁定。线程调用acquire()获取锁如果锁已被别的线程持有就阻塞等待拿到锁后执行临界区代码最后release()释放。import threading lock threading.Lock() counter 0 def worker(): global counter for _ in range(1000): lock.acquire() counter 1 lock.release() threads [threading.Thread(targetworker) for _ in range(4)] for t in threads: t.start() for t in threads: t.join() print(counter) # 4000这里必须注意几个细节。第一acquire()和release()要成对出现最稳妥的写法是用with lock:它保证即使临界区抛异常也会自动释放锁。第二临界区范围要尽量小——你把整个计算任务都罩上锁等于把多线程退化回了单线程。RLock是重入锁允许同一个线程多次acquire()而不死锁。典型场景是递归调用加锁一个函数内部获取了锁递归调用自身时又想获取同一把锁普通Lock会直接死锁RLock知道是同一个线程直接放行。rlock threading.RLock() def recursive(count): rlock.acquire() if count 0: recursive(count - 1) rlock.release()选锁的原则很简单临界区里没有嵌套获取同一把锁的需求就用Lock有递归、回调、层层调用场景直接用RLock别赌“永远不会重入”。2.2 Semaphore与BoundedSemaphore流量控制信号量维护一个计数器acquire()使计数器减1减到0时阻塞release()使计数器加1。它不限制“谁能进”只限制“同时最多进几个”非常适合控制连接数、限流等场景。比如你写一个爬虫目标网站最多允许5个并发连接于是建一个Semaphore(5)每个线程在发起请求前acquire()请求结束后release()。这样无论你后来开了20个线程同时飞出去的请求最多5个。import threading import time semaphore threading.Semaphore(5) def api_request(idx): semaphore.acquire() try: # 模拟发起请求 time.sleep(0.1) print(f请求 {idx} 完成) finally: semaphore.release()BoundedSemaphore是更安全的版本它在release()时检查计数器是否超过了初始值。普通信号量如果release()次数比acquire()多计数器会不断变大慢慢失去限流作用还不报错BoundedSemaphore会直接抛ValueError把这种“多放行”的bug暴露出来。我几乎永远用BoundedSemaphore贵不了几分钱安全多一份保障。2.3 Event一次性事件通知Event内部维护一个布尔标志三个常用方法set()把标志设为Truewait()阻塞直到标志为Trueclear()把标志重置为False。它适合“一个线程等待另一个线程完成某个动作再继续”的场景。我记忆很深的例子是程序启动时主线程需要等待后台线程完成一些初始化工作比如加载模型、建立连接池然后再开始处理任务。import threading init_event threading.Event() def init_worker(): # 模拟耗时初始化 time.sleep(2) print(后台初始化完成) init_event.set() def main_work(): print(等待初始化...) init_event.wait() print(开始主业务) threading.Thread(targetinit_worker).start() main_work()需要注意Event是“一次性开关”set()之后如果没有clear()后续wait()都会立即返回。如果业务需要反复通知并且要处理“通知来的时候接收方其实还没准备好”这种状态竞态Event就力不从心了这时候应该用Condition。2.4 Condition更复杂的条件等待Condition是锁和条件变量的组合核心是两个方法wait()会释放底层锁并阻塞直到被notify()或notify_all()唤醒notify()随机唤醒一个等待线程notify_all()唤醒全部。经典用法是生产者-消费者模型。消费者发现队列为空就wait()生产者往队列里塞数据后notify()唤醒消费者。注意一个铁律调用wait()和notify()前必须先获取锁否则直接抛RuntimeError。import threading condition threading.Condition() items [] def consumer(): with condition: while not items: condition.wait() item items.pop() print(f消费: {item}) def producer(): with condition: items.append(task-1) condition.notify() t threading.Thread(targetconsumer) t.start() producer() t.join()这里有一个关键细节消费者用while not items:而不是if not items:。因为wait()被唤醒后不保证条件一定满足——可能有多个消费者被唤醒其中一个先抢到锁把列表取空了后面醒来的线程再判断时条件已经不成立。while循环能防止这种“虚假唤醒”问题。2.5 Barrier凑齐线程再开工Barrier(parties)会让线程阻塞直到parties个线程都调用了wait()然后所有线程同时被释放。它适合“所有线程就绪后再同时开始”的同步场景。我做过一个并行数据集加载的任务8个线程各加载自己负责的数据分片加载完成后必须等所有分片都就绪才能进入下一阶段的聚合计算。用Barrier(8)正好import threading barrier threading.Barrier(8) def load_part(part_id): # 模拟加载耗时 time.sleep(part_id * 0.1) print(f分片 {part_id} 加载完成) barrier.wait() print(f分片 {part_id} 进入聚合阶段)Barrier还有个实用参数timeout在wait()里指定。如果等了超时还没凑齐BrokenBarrierError会被抛出方便你做异常处理而不是无限卡死。2.6 Queue最省心的线程安全方案queue.Queue内部自己用了锁和条件变量是线程安全的。生产者使用put()放入数据消费者使用get()取出数据两个方法都支持timeout参数。这是我最推荐的线程间数据交互方式。import queue import threading q queue.Queue(maxsize10) def producer(): for i in range(20): q.put(i) print(f生产: {i}) def consumer(): while True: item q.get() if item is None: # 哨兵值通知退出 break print(f消费: {item}) threading.Thread(targetproducer).start() threading.Thread(targetconsumer).start()用Queue的最大好处是你不用自己操心锁的获取和释放内部实现已经在并发环境下验证过很多年。上面例子里None作为哨兵值的退出方式是多线程编程里很常见的技巧比直接用stop_event更直观——数据流终止和线程退出的语义绑定在一起。3. 线程池与同步的搭配实战3.1 ThreadPoolExecutor常见误区concurrent.futures.ThreadPoolExecutor是内置线程池比手动管理线程省心得多。但很多新手拿着它当“无脑并行器”忽略了一个问题线程池里的任务也会共享状态也需要同步。线程池最常用的提交方式是submit(fn, *args)返回一个Future对象。调用future.result()会阻塞当前线程直到那个任务执行完毕并返回结果。如果任务内部抛了异常result()会原样把异常抛出来。from concurrent.futures import ThreadPoolExecutor, as_completed def fetch(url): # 模拟网络请求 return f结果: {url} with ThreadPoolExecutor(max_workers4) as pool: futures [pool.submit(fetch, furl-{i}) for i in range(10)] for future in as_completed(futures): print(future.result())as_completed会按任务完成顺序返回结果而不是按提交顺序。如果你需要严格的提交顺序处理结果直接对futures列表做循环调用result()就行。另外with语句会在退出时调用shutdown(waitTrue)自动等待所有任务结束这也是推荐写法。3.2 线程池内共享变量必须加锁线程池虽然管理了线程但你的业务代码里如果有共享变量血泪教训依然会重演。我写过一个统计任务线程池处理一批请求需要统计成功和失败的数量。如果直接在每个任务里写success_count 1跑完统计数据经常对不上。正确做法是在线程池外面创建一个Lock任务内部执行计数时加锁import threading from concurrent.futures import ThreadPoolExecutor lock threading.Lock() success_count 0 fail_count 0 def process_task(task): global success_count, fail_count try: result do_request(task) with lock: success_count 1 except Exception: with lock: fail_count 1 with ThreadPoolExecutor(max_workers8) as pool: list(pool.map(process_task, tasks)) print(f成功 {success_count}失败 {fail_count})这里再补充一个容易被忽略的点pool.map返回的是一个生成器只有在迭代时才会真正获取任务结果如果你不给生成器加list()或者不迭代它异常会被吞掉任务未必全部完成。所以写线程池任务时一定要确保Future被消费了。3.3 自定义线程池的阻塞队列怎么选如果你不用ThreadPoolExecutor而是自己手动实现一个线程池就要设计任务队列。Python内置的queue.Queue可以设置maxsize这就是有界队列不设置maxsize就是无界队列。经验法则生产消费速度差异大时用有界队列并搭配put(timeout)防止生产者无限阻塞。无界队列看似简单但任务堆积太多会吃掉大量内存再碰上消费端异常整个进程可能直接OOM。我在处理消息转发任务时就遇到过某个下游服务变慢任务队列从几百涨到几十万最后内存爆掉。后来给队列加上maxsize5000满的时候生产者等待并记录告警日志问题立刻可控。选择队列大小没有绝对公式我通常按“消费端峰值处理速度 × 可容忍的排队秒数”估算。比如消费端每秒能处理1000个任务我允许任务排队10秒那就是maxsize10000。这只是起点上线后还要根据内存和延迟指标微调。3.4 线程池等待所有任务完成的几种姿势线程池场景里有一个高频需求发出去一堆任务等待它们全部完成再做后续操作。姿势有几种不同场景选不同的。任务不多且不需要实时结果时用pool.shutdown(waitTrue)但shutdown之后这个线程池就不能再提交新任务了。如果还要继续提交改用concurrent.futures.wait(futures, return_whenALL_COMPLETED)它不关闭线程池只是把当前线程阻塞到所有任务完成。from concurrent.futures import ThreadPoolExecutor, wait with ThreadPoolExecutor(max_workers4) as pool: futures [pool.submit(job, i) for i in range(20)] done, not_done wait(futures, timeout60) print(f完成 {len(done)} 个任务超时未完成 {len(not_done)} 个)想同时拿到每个任务的结果并按完成顺序处理用as_completed配合future.result()。需要特别注意future.result()没有设置超时的话会无限等待我一般都会给它加个timeout参数万一任务卡住了也能及时暴露问题。4. 死锁与排查技巧实录4.1 死锁产生的四个条件死锁是线程同步里最讨厌的问题程序不报错、不崩溃就是卡住不动。产生死锁需要同时满足四个条件资源互斥一个资源同时只能被一个线程占用、持有并等待线程持有资源A等待资源B、不可剥夺资源不能被别人强行拿走、循环等待线程1持有A等B线程2持有B等A。代码层面的典型现场是这样的import threading lock_a threading.Lock() lock_b threading.Lock() def worker_1(): with lock_a: with lock_b: print(worker1 拿到两把锁) def worker_2(): with lock_b: with lock_a: print(worker2 拿到两把锁)两个线程按相反顺序加锁就会形成经典的“互相等待”。这跟两个人面对面过独木桥谁都不肯退让一样结果就是卡在桥上谁也走不了。4.2 我踩过的死锁现场光讲理论没用说一个我真实遇到的死锁案例。业务逻辑是线程A持有一个数据库连接锁等待另一个线程B把结果放入队列线程B在放入队列之前需要获取同一个数据库连接锁。两者各执一词都在等对方释放整个服务卡死。还有一次更隐蔽我在某个回调函数里调用了ThreadPoolExecutor实例的shutdown(waitTrue)而这个回调函数本身就在那个线程池的任务中执行——相当于自己等自己完成必然死锁。查了半个多小时才反应过来。线程池里的任务不能再调用shutdown(waitTrue)等待整个池子结束这是刻在骨子里的教训。这些案例说明一个道理死锁往往不是一眼能看出来的它藏在调用层级过深的协作逻辑里。排查时不能只盯着某个锁要看整个线程之间的依赖关系。4.3 死锁排查方法实战发现程序卡死时第一步是拿到所有线程的堆栈。我最常用的工具是py-spy一条命令搞定py-spy dump --pid pid它会把进程里每个线程当前执行到哪一行、正在获取什么锁都打出来。看到两个线程分别停在对方的acquire()上死锁原因基本就清楚了。import threading lock_a threading.Lock() lock_b threading.Lock() def worker_1(): with lock_a: print(t1 lock_a) lock_b.acquire() lock_b.release() def worker_2(): with lock_b: print(t2 lock_b) lock_a.acquire() lock_a.release()在开发环境里也可以提前给锁加超时避免无限阻塞if not lock_b.acquire(timeout5): print(获取锁超时可能存在死锁!) # 这里可以打印线程堆栈供分析注意Lock.acquire(timeout...)和RLock.acquire(timeout...)都支持超时参数但with lock:语法没法直接传超时。所以需要超时保护的场景还得老老实实用acquire(timeout...)配合try/finally。4.4 避免死锁的几个铁律我把多年总结的避坑规则列在这里每一条都是用线上事故换来的。固定加锁顺序所有线程获取多把锁时都遵循同一顺序比如先锁A再锁B永远不要反过来。这是最简单有效的死锁预防手段。缩小临界区锁里做的事情越少持锁时间越短死锁概率越低。可以把需要锁保护的代码压缩到极致业务处理放到锁外。优先用队列代替锁生产者-消费者模型用Queue天然不会死锁队列内部也有锁但实现已经处理好了顺序问题少写一个自定义锁就少一份风险。加锁超时而非无限等待生产环境尽量给acquire加超时一旦超时记录日志继续执行避免整个进程卡死。避免类锁的“自等”线程池任务里不要再调用shutdown(waitTrue)等自己结束回调函数也不要等自己所属线程池中的其他任务。4.5 常见问题速查表现象可能原因解决方案程序运行一段时间后卡死CPU占用低死锁用py-spy dump看线程栈检查加锁顺序计数器结果小于预期缺少同步保护给复合操作加锁或改用Queuefuture.result()一直阻塞任务未完成或任务内死锁加timeout参数用as_completed配合日志同一个锁acquire()两次直接卡住普通Lock不可重入改用RLockEvent.wait()总能立即返回事件状态被意外置位分析set()/clear()调用时机考虑用Condition信号量限流失效release次数超过acquire改用BoundedSemaphore线程池提交的任务没有全部执行生成器形式的map结果未被迭代显式迭代或转成list排查这类问题我个人的工作流永远是py-spy看栈 → 画出线程资源依赖图 → 检查加锁顺序 → 确认是否有“自等”。你不要指望肉眼读代码能发现所有问题工具和流程比记忆可靠得多。线程同步这件事做得多了我现在反而越来越“懒”——能交给Queue的绝不自己造锁能固定加锁顺序的绝不搞花活。你写的同步代码每多一行锁操作就多一分死锁风险。把状态变化先画清楚再决定用哪种原语比一上来就加锁靠谱得多。最后再分享一个小技巧开发环境里给所有自定义锁设置timeout5并打日志等于给你的并发代码装了个“烟雾报警器”线上出问题之前大概率能在测试阶段暴露出来。