Python threading.Condition与Barrier行为差异原因咨询
线程同步问题:Barrier与Condition的输出差异原因
需求背景
需要让线程在无限循环中执行工作(如IO操作),但每个周期需获得主线程的许可,且许可仅在所有线程完成上一周期后发放。尝试用threading.Barrier和threading.Condition实现,但两者输出差异显著,脚本2中线程输出完全没有交错,以下是具体分析。
脚本1:使用Barrier的实现
from threading import Barrier, Thread from time import sleep, time br = Barrier(3) store = [] def f1(): while True: br.wait() sleep(1) print("Calc part1") def f2(): while True: br.wait() sleep(1) print("Calc part2") Thread(target=f1).start() Thread(target=f2).start() for i in range(10): br.wait() print(f'end iter {i}') print(f'-------------')
预期输出(符合需求)
end iter 0 ------------- Calc part1 Calc part2 end iter 1 ------------- Calc part2 Calc part1 end iter 2 ------------- Calc part1 ...
脚本2:使用Condition的错误实现
from threading import Condition, Thread from time import sleep condition = False cv = Condition() def predicate(): return condition def f1(): for i in range(3): with cv: cv.wait_for(predicate) sleep(1) print("Calc part1") def f2(): for i in range(3): with cv: cv.wait_for(predicate) sleep(1) print("Calc part2") Thread(target=f1).start() Thread(target=f2).start() with cv: condition = True cv.notify_all()
实际输出(不符合预期)
Calc part1 Calc part1 Calc part1 Calc part2 Calc part2 Calc part2
差异原因分析
1. Barrier的工作逻辑(脚本1符合需求的原因)
Barrier(3)指定了需要3个线程(2个工作线程+主线程)同时到达wait()点才会继续执行:
- 每个周期内,主线程和两个工作线程都会调用
br.wait(),阻塞直到三者全部到达。 - 当所有线程都到达后,所有线程同时被唤醒,各自执行后续代码(主线程打印周期信息,工作线程sleep后打印任务内容)。
- 工作线程被唤醒后的执行顺序由操作系统调度决定,因此输出会随机交错,符合需求。
2. Condition的错误使用(脚本2输出无交错的核心原因)
脚本2的问题出在锁的持有逻辑和同步逻辑缺失:
- 锁长期持有:工作线程的
with cv会持有Condition关联的互斥锁,cv.wait_for(predicate)在满足条件返回后,锁仍然被线程持有。也就是说,线程在执行sleep(1)和print时,其他线程根本无法进入with cv块。第一个被唤醒的线程(比如f1)会一直握着锁,完成3次循环的所有操作后才释放锁,此时f2线程才能开始执行自己的循环,导致输出完全串行。 - 同步逻辑不符合周期需求:主线程仅设置一次
condition=True并通知所有线程,之后condition一直为True,工作线程的wait_for(predicate)会直接通过,完全没有实现“每个周期需等待所有线程完成后,主线程再发放许可”的逻辑。
修正Condition实现的思路
要让Condition实现类似Barrier的周期同步,需要满足以下几点:
- 工作线程在执行任务前等待主线程的许可,执行任务时释放锁(避免阻塞其他线程)。
- 每个周期结束后,工作线程通知主线程自己已完成,主线程等待所有线程完成后,再重置状态并发放下一个周期的许可。
修正后的示例代码:
from threading import Condition, Thread from time import sleep condition = False cv = Condition() completed_threads = 0 total_work_threads = 2 def f1(): global completed_threads while True: # 等待主线程的周期许可 with cv: cv.wait(lambda: condition) # 释放锁后执行耗时任务,不阻塞其他线程 sleep(1) print("Calc part1") # 通知主线程当前线程已完成周期任务 with cv: completed_threads += 1 cv.notify() def f2(): global completed_threads while True: with cv: cv.wait(lambda: condition) sleep(1) print("Calc part2") with cv: completed_threads += 1 cv.notify() Thread(target=f1).start() Thread(target=f2).start() # 主线程控制周期同步 for i in range(10): with cv: # 等待所有工作线程完成上一周期 while completed_threads < total_work_threads: cv.wait() # 重置状态,发放下周期许可 completed_threads = 0 condition = True cv.notify_all() print(f'end iter {i}') print(f'-------------') # 短暂等待让工作线程先启动任务,避免主线程快速进入下一轮 sleep(0.1)
内容的提问来源于stack exchange,提问作者eugeny
相关产品推荐
相关产品推荐

