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

如何实现新线程优先的线程调度:旧线程暂停恢复及代码问题解决

解决方案:高优先级新线程的线程调度实现

你的代码核心问题在于未使用线程安全的方式访问类变量newest_thread,多线程环境下读写该变量会产生竞态条件,导致判断逻辑失效;同时watch线程的循环没有休眠,会占用大量CPU资源,且逻辑判断存在漏洞。以下是修复并优化后的实现:

import threading
import time


class DialogueThread:
    # 类级变量:存储当前最高优先级的线程(最新创建的)
    _newest_thread = None
    # 类级锁:保护_newest_thread的读写,避免竞态条件
    _lock = threading.Lock()

    @classmethod
    def set_newest_thread(cls, thread):
        with cls._lock:
            cls._newest_thread = thread

    @classmethod
    def get_newest_thread(cls):
        with cls._lock:
            return cls._newest_thread

    def __init__(self, task):
        self.task = task
        # 用于暂停/恢复当前工作线程的事件
        self.pause_event = threading.Event()
        self.pause_event.set()  # 初始状态为允许运行
        # 工作线程:传入当前DialogueThread实例,让任务能获取新线程引用
        self.thread = threading.Thread(target=self.task, args=[self])
        self.thread.start()
        # 启动监控线程
        self.watch_thread = threading.Thread(target=self._watch_thread)
        self.watch_thread.start()
        # 将自身设为最新线程(最高优先级)
        DialogueThread.set_newest_thread(self.thread)

    def _watch_thread(self):
        while self.thread.is_alive():
            # 线程安全地获取当前最新线程
            current_newest = DialogueThread.get_newest_thread()
            # 如果当前线程不是最高优先级,暂停运行
            if current_newest is not None and current_newest != self.thread:
                self.pause_event.clear()
                # 等待最新线程结束(带超时,避免阻塞)
                current_newest.join(timeout=0.1)
            else:
                # 无更高优先级线程,恢复运行
                self.pause_event.set()
            # 短休眠降低CPU占用
            time.sleep(0.1)
        # 工作线程结束后,若自身是最新线程则清空
        with DialogueThread._lock:
            if DialogueThread._newest_thread == self.thread:
                DialogueThread._newest_thread = None

    def get_newest_thread_ref(self):
        # 提供给任务线程获取最新线程引用的方法
        return DialogueThread.get_newest_thread()


def thread1(dialogue_thread: DialogueThread):
    count = 0
    while count < 10:
        dialogue_thread.pause_event.wait()

        # 获取当前最新线程的引用
        newest_thread = dialogue_thread.get_newest_thread_ref()
        if newest_thread and newest_thread != dialogue_thread.thread:
            print(f"thread1 detected new thread ID: {newest_thread.ident}")
        
        print(f"thread 1 is working: {count}")
        count += 1
        time.sleep(1)
    print("thread1 is finished")


def thread2(dialogue_thread: DialogueThread):
    count = 0
    while count < 5:
        dialogue_thread.pause_event.wait()

        # 获取当前最新线程的引用
        newest_thread = dialogue_thread.get_newest_thread_ref()
        print(f"thread2 detected newest thread ID: {newest_thread.ident if newest_thread else None}")
        
        print(f"thread 2 is working: {count}")
        count += 1
        time.sleep(1)
    print("thread2 is finished")


if __name__ == "__main__":
    thread_1 = DialogueThread(thread1)
    time.sleep(5)
    thread_2 = DialogueThread(thread2)

关键改动说明

  • 线程安全的类变量管理:新增_lock类级锁,所有对_newest_thread的读写操作都通过锁保护,彻底解决多线程竞态问题,这是原代码判断失效的核心原因。
  • 任务线程传递DialogueThread实例:不再仅传递暂停事件,而是传入整个DialogueThread对象,让任务线程可以调用方法安全获取最新线程的引用,满足你“每个线程能获取新线程引用”的需求。
  • 优化监控线程逻辑:
    • 每次先线程安全地获取最新线程引用,再判断是否需要暂停当前线程。
    • 检测到更高优先级线程时,通过join(timeout=0.1)等待其结束,减少无效轮询次数。
    • 加入time.sleep(0.1)降低CPU占用率。
  • 新增线程引用获取方法:get_newest_thread_ref方法让任务线程可以便捷、安全地获取当前最高优先级线程的引用。

运行代码后,thread1运行5秒后thread2启动,此时thread1会立即暂停,直到thread2执行完成后,thread1才会恢复执行剩余循环;同时每个线程都能打印出最新线程的ID,验证可以正确获取新线程引用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 03:05:34