【问题标题】:Python: nonblocking read from stdout of threaded subprocessPython:从线程子进程的标准输出进行非阻塞读取
【发布时间】:2010-03-18 14:54:37
【问题描述】:

我有一个脚本 (worker.py),它以表单的形式打印无缓冲的输出...

1
2
3
.
.
.
n

其中 n 是此脚本中的循环将进行的某个恒定迭代次数。在另一个脚本(service_controller.py)中,我启动了许多线程,每个线程使用 subprocess.Popen(stdout=subprocess.PIPE, ...);现在,在我的主线程 (service_controller.py) 中,我想读取每个线程的 worker.py 子进程的输出,并使用它来计算完成前剩余时间的估计值。

我拥有从 worker.py 读取标准输出并确定最后打印的数字的所有逻辑。问题是我无法弄清楚如何以非阻塞方式执行此操作。如果我读取一个恒定的 bufsize,那么每次读取最终都会等待来自每个工作人员的相同数据。我尝试了很多方法,包括使用 fcntl、select + os.read 等。我最好的选择是什么?如果需要,我可以发布我的来源,但我认为解释很好地描述了问题。

感谢您在此提供的任何帮助。

编辑
添加示例代码

我有一个启动子进程的工人。

class WorkerThread(threading.Thread):
    def __init__(self):
        self.completed = 0
        self.process = None
        self.lock = threading.RLock()
        threading.Thread.__init__(self)

    def run(self):
        cmd = ["/path/to/script", "arg1", "arg2"]
        self.process = subprocess.Popen(cmd, stdout=subprocess.PIPE, bufsize=1, shell=False)
        #flags = fcntl.fcntl(self.process.stdout, fcntl.F_GETFL)
        #fcntl.fcntl(self.process.stdout.fileno(), fcntl.F_SETFL, flags | os.O_NONBLOCK)

    def get_completed(self):
        self.lock.acquire();
        fd = select.select([self.process.stdout.fileno()], [], [], 5)[0]
        if fd:
            self.data += os.read(fd, 1)
            try:
                self.completed = int(self.data.split("\n")[-2])
            except IndexError:
                pass
        self.lock.release()
        return self.completed

然后我有一个 ThreadManager。

class ThreadManager():
    def __init__(self):
        self.pool = []
        self.running = []
        self.lock = threading.Lock()

    def clean_pool(self, pool):
        for worker in [x for x in pool is not x.isAlive()]:
            worker.join()
            pool.remove(worker)
            del worker
        return pool

    def run(self, concurrent=5):
        while len(self.running) + len(self.pool) > 0:
            self.clean_pool(self.running)
            n = min(max(concurrent - len(self.running), 0), len(self.pool))
            if n > 0:
                for worker in self.pool[0:n]:
                    worker.start()
                self.running.extend(self.pool[0:n])
                del self.pool[0:n]
            time.sleep(.01)
         for worker in self.running + self.pool:
             worker.join()

以及一些运行它的代码。

threadManager = ThreadManager()
for i in xrange(0, 5):
    threadManager.pool.append(WorkerThread())
threadManager.run()

我已经删除了其他代码的日志,希望能找出问题所在。

【问题讨论】:

  • 您使用的是 Linux 还是其他 Unix?如果是这样, select + os.read 1 byte 应该可以正常工作——你能告诉我们你在这条线上的代码以及它给你带来了什么错误或不当行为吗?
  • 这实际上是在windoze上运行的,用于开发将在Fedora或OS X上用于生产。

标签: python multithreading subprocess nonblocking


【解决方案1】:

不要让你的 service_controller 被 i/o 访问阻塞,只有线程循环应该读取它自己的受控进程输出。

然后,您可以在线程对象中使用方法来控制进程以获取最后轮询的输出。

当然,在这种情况下不要忘记使用某种锁定机制来保护缓冲区,线程将使用该缓冲区来填充它以及控制器调用的方法来获取它。

【讨论】:

  • 我与您的建议相差甚远吗?我有线程对象控制进程以获取最后轮询输出...
  • 您的 get_completed 方法仅填充 self.completed ,我建议将其重命名为 update_completed。然后添加一个返回 self.completed 的 get_completed 方法(添加一个 threading.RLock 来保护对它的访问)。然后在您的线程管理器中,您可以定期在您的工作人员上调用 get_completed。
  • get_completed 方法实际上应该返回 self.completed (我在重新输入时错误地省略了它)。我已经在读数周围添加了 RLock,但我仍然遇到同样的问题。
  • 我无法理解这一点。也许今晚睡一觉会让我有我正在寻找的突破。
【解决方案2】:

您的方法 WorkerThread.run() 启动子进程,然后立即终止。 Run() 需要执行轮询并更新 WorkerThread.completed 直到子进程完成。

【讨论】:

    猜你喜欢
    • 2014-03-29
    • 1970-01-01
    • 2017-12-30
    • 2016-07-28
    • 2011-09-18
    • 1970-01-01
    • 2018-06-23
    • 2012-01-18
    • 1970-01-01
    相关资源
    最近更新 更多