【问题标题】:Streaming wrapper around program that writes to multiple output files围绕写入多个输出文件的程序的流式包装器
【发布时间】:2023-03-13 07:04:01
【问题描述】:

有一个程序(我无法修改)创建两个输出文件。我正在尝试编写一个 Python 包装器来调用该程序,同时读取两个输出流,组合输出,并打印到标准输出(以促进流式传输)。我怎样才能做到这一点而不会死锁?下面的概念证明工作正常,但是当我将此方法应用于实际程序时,它会死锁。


概念证明:这是一个虚拟程序bogus.py,它创建两个输出文件,就像我要包装的程序一样。

#!/usr/bin/env python
from __future__ import print_function
import sys
with open(sys.argv[1], 'w') as f1, open(sys.argv[2], 'w') as f2:
    for i in range(1000):
        if i % 2 == 0:
            print(i, file=f1)
        else:
            print(i, file=f2)

这里是 Python 包装器,它调用程序并组合它的两个输出(每次交错 4 行)。

#!/usr/bin/env python
from __future__ import print_function
from contextlib import contextmanager
import os
import shutil
import subprocess
import tempfile

@contextmanager
def named_pipe():
    """
    Create a temporary named pipe.

    Stolen shamelessly from StackOverflow:
    http://stackoverflow.com/a/28840955/459780
    """
    dirname = tempfile.mkdtemp()
    try:
        path = os.path.join(dirname, 'named_pipe')
        os.mkfifo(path)
        yield path
    finally:
        shutil.rmtree(dirname)

with named_pipe() as f1, named_pipe() as f2:
    cmd = ['./bogus.py', f1, f2]
    child = subprocess.Popen(cmd)
    with open(f1, 'r') as in1, open(f2, 'r') as in2:
        buff = list()
        for i, lines in enumerate(zip(in1, in2)):
            line1 = lines[0].strip()
            line2 = lines[1].strip()
            print(line1)
            buff.append(line2)
            if len(buff) == 4:
                for line in buff:
                    print(line)

【问题讨论】:

  • 您是否尝试过显而易见的:subprocess.check_call(['/program', '-', '-']),如果program 不理解'-',则通过'/dev/stdout''/dev/fd/1' 或使用单个命名管道而不是'-'取决于系统。注意:如果子进程同时在自己的线程中写入文件,则输出可能会在一行中间交错。
  • 是的,使用/dev/stdout作为输出文件名时,输出的顺序有问题。它们确实需要被视为单独的流。
  • 这里的“订单”是什么意思?你能举个例子吗? 1- 不同文件的写入顺序未定义(子进程外部没有之前或之后,例如,如果子进程写入:write(fd1, 'a'); write(fd2, 'b') 并且您正在从与fd1fd2 对应的文件中读取那么您无法知道'a''b' 的优先顺序。2-如果您在谈论顺序,则可能是在谈论子进程中内部缓冲的表现...
  • ..[继续] 问题可能是标准输出缓冲。查看传递/dev/stderr/dev/tty 是否有助于说服孩子将其输出线缓冲到文件中。如果您无法控制孩子的内部缓冲,当您可能从一个文件中看到一大块然后从另一个文件中看到一大块时,等等
  • 我看到一个文件的大块,然后是另一个文件的大块,无论我是写入 stdout、stderr 还是 tty。

标签: python subprocess deadlock pipeline


【解决方案1】:

我看到的是一个文件的大块,然后是另一个文件的大块,无论我是写入 stdout、stderr 还是 tty。

如果您不能让孩子对文件使用行缓冲,那么一个简单的解决方案 在输出可用时进程仍在运行时从输出文件中读取完整的交错行 是使用线程:

#!/usr/bin/env python2
from subprocess import Popen
from threading import Thread
from Queue import Queue

def readlines(path, queue):
    try:
        with open(path) as pipe:
            for line in iter(pipe.readline, ''):
                queue.put(line)
    finally:
        queue.put(None)

with named_pipes(n=2) as paths:
    child = Popen(['python', 'child.py'] + paths)
    queue = Queue()
    for path in paths:
        Thread(target=readlines, args=[path, queue]).start()
    for _ in paths:
        for line in iter(queue.get, None):
            print line.rstrip('\n')

named_pipes(n) is defined here.

pipe.readline() 在 Python 2 上因非阻塞管道而损坏,这就是此处使用线程的原因。


从一个文件中打印一行,然后从另一个文件中打印一行:

with named_pipes(n=2) as paths:
    child = Popen(['python', 'child.py'] + paths)
    queues = [Queue() for _ in paths]
    for path, queue in zip(paths, queues):
        Thread(target=readlines, args=[path, queue]).start()
    while queues:
        for q in queues:
            line = q.get()
            if line is None:  # EOF
                queues.remove(q)
            else:
                print line.rstrip('\n')

如果child.py 向一个文件写入的行多于另一个文件,则差异将保存在内存中,因此queues 中的各个队列可能会无限增长,直到填满所有内存。您可以设置队列中的最大项目数,但您必须将超时传递给q.get(),否则代码可能会死锁。


如果您需要从一个输出文件准确打印 4 行,然后从另一个输出文件准确打印 4 行,等等,那么您可以稍微修改给定的代码示例:

    while queues:
        # print 4 lines from one queue followed by 4 lines from another queue
        for q in queues:
            for _ in range(4):
                line = q.get()
                if line is None:  # EOF
                    queues.remove(q)
                    break
                else:
                    print line.rstrip('\n')

它不会死锁,但如果您的子进程将太多数据写入一个文件而没有向另一个文件写入足够多的数据,它可能会占用所有内存(只有差异保留在内存中 - 如果文件相对相等;程序支持任意大的输出文件)。

【讨论】:

  • 太棒了,这解决了死锁问题。但是,输出的顺序仍然是一个问题。它需要是第一个文件的 4 行,然后是下一个文件的 4 行,依此类推。我仍然一次从每个文件中获取大块。
  • 想法:将for path in paths: 替换为for i, path in enumerate(paths):,然后将i 作为参数传递给readlines。这是否允许我将数据累积到两个单独的队列中并一次从每个队列中弹出 4 行?
  • 这很复杂,因为这通常会处理非常大量的数据,所以我们不能将它们全部累积到内存中,然后在最后处理。
  • @DanielStandage 如果文件永远不会太不同步,那么修改代码以交错每一行(第一个文件中的一行,下一行来自另一个文件等)是很简单的—尽管如果允许文件不同步(例如,如果第一个文件比另一个文件大 10 倍,那么修改后的代码将死锁(对于足够大的文件),而我的答案中的代码有效)
【解决方案2】:

Popen 只产生进程。您必须执行类似child.communicate() 之类的操作才能实际与之交互并获取其输出。

另外,我认为您需要在开始该过程之前open 管道进行读取。

【讨论】:

  • 如果我理解正确,communicate() 只会返回标准输入/标准输出。我想截取程序写入两个输出文件的数据。在我的概念验证中,在调用 Popen 之前打开管道进行读取会导致死锁。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-02-22
  • 2013-11-12
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多