多线程下实现do方法与指定方法的同步执行问题求助
多线程同步问题解决方案:do方法与指定方法的互斥等待
需求概述
- Bar实例的方法可在任意线程中被任意次数、任意顺序调用
- 非指定方法启动后可独立并行执行
- 一旦do方法启动,需等待所有指定方法(含其他do)执行完毕;do执行期间,其他指定方法需等待当前do完成才能继续;do执行完毕后,其他指定方法才可恢复执行
- 仅
FUNC_SYNCHRONISE中指定的方法需要同步,其他方法不受限制 - 必须保留
Bar(Foo)的继承关系,且不能修改Foo类
现有代码问题分析
原实现的核心问题:
- 使用
threading.Event无法正确处理多个do的排队逻辑,多个do线程会同时通过wait()进入执行流程 queue.Queue的get(name)用法错误,Queue的get()是取出队首元素,并非按值查找- 没有保证状态检查与修改的原子性,导致并发场景下出现竞态条件
修正后的代码实现
import threading from concurrent.futures import ThreadPoolExecutor import time import logging logging.basicConfig( level=logging.INFO, format="%(message)s", ) class Foo(object): def func_1(self): log = logging.getLogger(__name__) log.info("start func_1") time.sleep(2) log.info("end func_1") def func_2(self): log = logging.getLogger(__name__) log.info("start func_2") time.sleep(2) log.info("end func_2") def func_3(self): log = logging.getLogger(__name__) log.info("start func_3") time.sleep(2) log.info("end func_3") def func_4_free(self): log = logging.getLogger(__name__) log.info("start func_4_free") time.sleep(2) log.info("end func_4_free") def do(self): log = logging.getLogger(__name__) log.info("start do") time.sleep(1) log.info("This is do") time.sleep(1) log.info("end do") class Bar(Foo): # 指定需要同步的方法:key是触发独占的方法,value是需要同步的方法集合 FUNC_SYNCHRONISE = {"do": ["func_1", "func_2", "func_3", "do"]} def __init__(self): super().__init__() # 保护共享状态的锁 self._lock = threading.Lock() # 用于线程等待/通知的条件变量 self._cond = threading.Condition(self._lock) # 当前正在执行的同步方法数量(不含等待中的) self._active_sync_funcs = 0 # 是否有do方法正在执行 self._do_running = False def __getattribute__(self, name): # 避免递归调用__getattribute__ if name in ["_lock", "_cond", "_active_sync_funcs", "_do_running", "FUNC_SYNCHRONISE"]: return super().__getattribute__(name) attr = super().__getattribute__(name) sync_methods = self.FUNC_SYNCHRONISE["do"] # 如果是需要同步的方法 if name in sync_methods: def wrapped(*args, **kwargs): log = logging.getLogger(__name__) with self._cond: if name == "do": # do方法需要等待所有同步方法结束,且没有其他do在运行 log.info(f"do waiting for all sync functions to finish...") while self._active_sync_funcs > 0 or self._do_running: self._cond.wait() # 标记当前do开始执行 self._do_running = True else: # 普通同步方法需要等待do执行结束 while self._do_running: self._cond.wait() # 标记当前同步方法开始执行 self._active_sync_funcs += 1 try: # 执行原方法逻辑 attr(*args, **kwargs) finally: with self._cond: if name == "do": # do执行完毕,标记状态并通知所有等待线程 self._do_running = False log.info(f"do finished, waking up waiting threads") else: # 普通同步方法执行完毕,减少计数 self._active_sync_funcs -= 1 # 通知所有等待的线程状态变化 self._cond.notify_all() return wrapped else: # 非同步方法直接返回原方法 return attr with ThreadPoolExecutor(max_workers=None) as executor: obj = Bar() # 测试用例:任意顺序、任意次数调用方法 test_calls = [ obj.func_1, obj.func_3, obj.func_4_free, obj.do, obj.func_2, obj.func_1, obj.do, obj.func_4_free, obj.do, obj.func_3, obj.func_2, obj.func_1 ] for func in test_calls: executor.submit(func)
核心逻辑解释
- 同步原语选择:使用
threading.Condition配合内置的Lock,保证状态检查与修改的原子性,同时支持线程的等待与通知 - 状态维护:
_active_sync_funcs:统计当前正在执行的非do同步方法数量_do_running:标记是否有do方法正在执行,确保同一时间只有一个do在运行
- 方法包装逻辑:
- do方法:获取锁后,等待所有同步方法(包括其他do)执行完毕,然后标记
_do_running为True,执行原方法;执行完毕后重置状态并通知所有等待线程 - 普通同步方法:获取锁后,等待
_do_running为False,然后增加_active_sync_funcs计数,执行原方法;执行完毕后减少计数并通知等待线程 - 非同步方法:直接返回原方法,不受同步逻辑限制
- do方法:获取锁后,等待所有同步方法(包括其他do)执行完毕,然后标记
内容的提问来源于stack exchange,提问作者Anton
相关产品推荐
相关产品推荐

