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

多层QThread架构下QProcess的readyReadStandardOutput无法正常输出

问题解决

问题根源

  1. Task线程阻塞导致信号无法处理:Task类的run方法里调用了waitForFinished(-1),这会完全阻塞线程,Qt事件循环无法处理QProcess发出的started和readyReadStandardOutput信号,因此对应的槽函数无法执行。
  2. AllocateParallel无事件循环导致调度低效:AllocateParallel的run方法用轮询方式等待任务结束,没有启动Qt事件循环,不仅效率低,还可能影响跨线程信号的处理。

修复方案

1. 修改Task类

移除阻塞式的waitForFinished,改为启动线程事件循环,通过QProcess的finished信号触发线程退出,确保信号能被正常分发和处理。

2. 修改AllocateParallel类

将轮询逻辑改为信号驱动:当一个任务结束时,自动尝试启动新任务,同时启动自身的事件循环来处理跨线程信号,提升调度效率。

修改后的完整代码

from PyQt5.QtCore import *
from PyQt5.QtWidgets import QApplication
from collections import deque
import sys

class AllocateParalle(QThread):
    def __init__(self, queue_test):
        super().__init__()
        self.queue_test = queue_test
        self.activate = 0
        self.max_activate = 1
        self.time = 0
        self.mutex = QMutex()
        self.list_thread = []

    def run(self):
        # 初始化启动最大并发数的任务
        for _ in range(self.max_activate):
            if self.queue_test:
                self.start_next_task()
        # 启动线程事件循环,处理信号
        self.exec()

    def start_next_task(self):
        self.mutex.lock()
        if self.activate < self.max_activate and self.queue_test:
            self.activate += 1
            self.time += 1
            cmd = self.queue_test.popleft()
            self.mutex.unlock()
            
            thread = Task(cmd, self.time)
            thread.finished.connect(self.finished_process)
            self.list_thread.append(thread)
            thread.start()
        else:
            self.mutex.unlock()

    def finished_process(self):
        self.mutex.lock()
        self.activate -= 1
        self.mutex.unlock()
        # 任务结束后尝试启动下一个任务
        self.start_next_task()

class Task(QThread):
    def __init__(self, cmd, idx):
        super().__init__()
        self.lsf_task_process = None
        self.case_cmd = cmd
        self.idx = idx

    def run(self):
        self.lsf_task_process = QProcess()
        self.lsf_task_process.setProcessChannelMode(QProcess.MergedChannels)

        env = QProcessEnvironment.systemEnvironment()
        env.insert('PATH', '/usr/bin:'+env.value('PATH'))
        self.lsf_task_process.setProcessEnvironment(env)

        # 连接信号到槽函数
        self.lsf_task_process.readyReadStandardOutput.connect(self.output)
        self.lsf_task_process.started.connect(self.start_output)
        self.lsf_task_process.finished.connect(self.process_finished)
        # 进程结束后退出线程事件循环
        self.lsf_task_process.finished.connect(self.quit)

        self.lsf_task_process.start('bash', ['-c', self.case_cmd])
        
        # 启动线程事件循环,处理信号
        self.exec()

    def start_output(self):
        print('start remote')

    def output(self):
        while self.lsf_task_process.canReadLine():
            line_data = self.lsf_task_process.readLine().data().decode('utf-8', errors='replace').strip()
            print(f'[{self.idx}] {line_data}')

    def process_finished(self, exitcode):
        if exitcode != 0:
            print(f'[{self.idx}] error')
        else:
            print(f'[{self.idx}] success')

if __name__ =="__main__":
    app = QApplication(sys.argv)
    queue = deque()
    cmd = "for i in {1..5}; do echo 'Hello, world!'; sleep 1; done"
    for i in range(3):
        queue.appendleft(cmd)
    thread = AllocateParalle(queue)
    # thread = Task(cmd, 0)
    thread.start()
    sys.exit(app.exec())

关键修改点

  • Task类:移除waitForFinished(-1),添加self.exec()启动事件循环,通过QProcess的finished信号连接self.quit()来结束线程。
  • AllocateParallel类:新增start_next_task方法处理任务启动逻辑,用finished信号触发下一个任务启动,同时启动自身事件循环self.exec()。
  • 用Qt的QMutex替代threading.Lock,更符合Qt线程模型的规范。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 13:15:53