QApplication中线程间通信规范实现与多线程同步问询
问题解答与代码优化建议
一、ThreadWorker主线程接收信号的问题与规范方案
问题分析
从日志可见ThreadWorker的槽函数(set_mode、start_it、stop_it)在主线程执行,这会引发线程安全隐患:
- ThreadWorker的核心工作逻辑在子线程(
thread1)运行,会读取self.mode、self.paused等状态变量; - 修改这些变量的槽函数却在主线程执行,无同步机制时,子线程可能读取到中间状态的脏数据,导致暂停指令延迟生效、模式切换混乱等问题。
规范实现方式
你想用全局统一Signals类的思路是合理的,只需调整信号连接时机,确保槽函数在ThreadWorker所在的子线程执行:
- 延迟信号连接:在ThreadWorker对象
moveToThread之后再绑定信号。修改MainWindow的线程初始化代码:
# 原代码 self.thread_worker = ThreadWorker(self.signals) self.thread_worker.moveToThread(self.thread1) self.signals.refreshGraphs.connect(self.plot_data) self.thread1.started.connect(self.thread_worker.start_com)
改为:
self.thread_worker = ThreadWorker() self.thread_worker.moveToThread(self.thread1) # 此时ThreadWorker已在子线程,槽函数会自动在子线程执行 self.signals.changeMode.connect(self.thread_worker.set_mode) self.signals.start.connect(self.thread_worker.start_it) self.signals.stop.connect(self.thread_worker.stop_it) # 其他连接保持不变 self.signals.refreshGraphs.connect(self.plot_data) self.thread1.started.connect(self.thread_worker.start_com)
同时修改ThreadWorker的__init__,移除内部的信号连接:
def __init__(self): log.info("Inside thread init") super(ThreadWorker, self).__init__() self.mode = "running" self.paused = True self._count = 0 self.refresh_rate = 1 self.average = 0 self.current_value = 0 self.ratio = 0 self.buffer = deque([np.nan], 100) self.cal_vals = [] self.reader = ValueReader(Signals()) # 使用单例Signals获取实例
- 锁保护状态变量:如果必须在主线程修改子线程状态,给
self.mode、self.paused加锁:
def __init__(self): # ... 其他初始化代码 self.lock = threading.Lock() def start_it(self): log.info("Inside of the start_it method of ThreadWorker.") with self.lock: self.paused = False def run_data_thread(self): while True: with self.lock: paused = self.paused mode = self.mode if paused: break # 后续逻辑使用paused和mode变量
二、第三个线程基于第二个线程信息的通信方案
要避免基于旧数据决策,核心是保证数据实时性与线程安全传递,可采用以下方案:
线程安全队列传递:
- 第二个线程(数据采集线程)将最新数据放入
queue.Queue(Python内置线程安全队列),设置队列最大长度为1,新数据自动覆盖旧数据; - 第三个线程从队列取数据时,直接获取最新值。
示例:
import queue # 全局或通过单例Signals传递队列 data_queue = queue.Queue(maxsize=1) # 第二个线程发送数据 def run_data_thread(self): # ... 获取new_data后 try: data_queue.put_nowait(new_data) except queue.Full: data_queue.get_nowait() data_queue.put_nowait(new_data) # 第三个线程处理数据 def third_thread_run(self): while not self.interrupted: if not data_queue.empty(): latest_data = data_queue.get() # 基于latest_data执行决策逻辑 time.sleep(0.01)- 第二个线程(数据采集线程)将最新数据放入
带版本号的共享数据:
- 维护一个带版本号的共享结构,比如
{"version": int, "data": any}; - 第二个线程更新数据时递增版本号,第三个线程只处理版本号高于上次记录的数据;
- 用锁保护数据的读写操作。
- 维护一个带版本号的共享结构,比如
PyQt信号传递最新数据:
- 在
Signals类中新增信号newDataReady = pyqtSignal(object); - 第二个线程采集到新数据时发出该信号;
- 第三个线程连接该信号,槽函数中直接处理最新数据(若担心信号堆积,可在槽函数中维护一个“最新数据”变量,每次信号触发时更新,处理时取该变量值)。
- 在
三、其他代码问题与优化建议
嵌套循环冗余:
ThreadWorker的start_com和run_data_thread都是while循环,暂停后会频繁进入空循环浪费资源,合并逻辑:def start_com(self): log.info("Inside of the start_com method of ThreadWorker.") while not self.thread().isInterruptionRequested(): with self.lock: paused = self.paused if not paused: self.current_value = self.reader.read_value() with self.lock: mode = self.mode if mode == "running": self.ratio = self.current_value self.buffer.append(self.current_value) self.signals.refreshGraphs.emit(self.ratio, self._count) elif mode == "calibrate": self.ratio = self.current_value self.buffer.append(self.current_value) self.signals.refreshGraphs.emit(self.ratio, self._count) self.calibrate(self.current_value) self._count += 1 if self._count == 100: self._count = 0 time.sleep(1 / self.refresh_rate) else: time.sleep(0.1) # 暂停时低频率轮询cal_vals初始化错误:
ThreadWorker的calibrate方法中,self.cal_vals = [[], [], []]会将列表改为二维数组,后续append(ratio)会出错,应改为self.cal_vals = []。跨线程状态访问优化:
MainWindow的_calibrate方法直接访问self.thread_worker.paused属于跨线程操作,应通过信号让ThreadWorker返回状态,或加锁读取。Signals改为单例:
将Signals改为单例模式,避免多次创建,方便其他类直接获取实例:class Signals(QObject, metaclass=Singleton): changeMode = pyqtSignal(str) refreshGraphs = pyqtSignal(float, int) start = pyqtSignal() stop = pyqtSignal(bool) calibrationVal = pyqtSignal(float)
内容的提问来源于stack exchange,提问作者Amanda.py
相关产品推荐
相关产品推荐

