You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

多线程下实现do方法与指定方法的同步执行问题求助

多线程同步问题解决方案:do方法与指定方法的互斥等待

需求概述

  • Bar实例的方法可在任意线程中被任意次数、任意顺序调用
  • 非指定方法启动后可独立并行执行
  • 一旦do方法启动,需等待所有指定方法(含其他do)执行完毕;do执行期间,其他指定方法需等待当前do完成才能继续;do执行完毕后,其他指定方法才可恢复执行
  • 仅FUNC_SYNCHRONISE中指定的方法需要同步,其他方法不受限制
  • 必须保留Bar(Foo)的继承关系,且不能修改Foo类

现有代码问题分析

原实现的核心问题:

  1. 使用threading.Event无法正确处理多个do的排队逻辑,多个do线程会同时通过wait()进入执行流程
  2. queue.Queue的get(name)用法错误,Queue的get()是取出队首元素,并非按值查找
  3. 没有保证状态检查与修改的原子性,导致并发场景下出现竞态条件

修正后的代码实现

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)

核心逻辑解释

  1. 同步原语选择:使用threading.Condition配合内置的Lock,保证状态检查与修改的原子性,同时支持线程的等待与通知
  2. 状态维护:
    • _active_sync_funcs:统计当前正在执行的非do同步方法数量
    • _do_running:标记是否有do方法正在执行,确保同一时间只有一个do在运行
  3. 方法包装逻辑:
    • do方法:获取锁后,等待所有同步方法(包括其他do)执行完毕,然后标记_do_running为True,执行原方法;执行完毕后重置状态并通知所有等待线程
    • 普通同步方法:获取锁后,等待_do_running为False,然后增加_active_sync_funcs计数,执行原方法;执行完毕后减少计数并通知等待线程
    • 非同步方法:直接返回原方法,不受同步逻辑限制

内容的提问来源于stack exchange,提问作者Anton

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.17 06:45:16