【问题标题】:allowing multiple inputs to python subprocess允许对 python 子进程进行多个输入
【发布时间】:2015-07-23 14:07:58
【问题描述】:

我有一个与几年前问过的一个几乎相同的问题:Python subprocess with two inputs 收到了一个答案但没有实施。我希望这篇转发可以帮助我和其他人澄清问题。

如上所述,我想使用 subprocess 来包装一个命令行工具,该工具需要多个输入。特别是,我想避免将输入文件写入磁盘,而是宁愿使用例如命名管道,正如上面提到的。那应该是“学习如何”,因为我承认我以前从未尝试过使用命名管道。我将进一步说明,我目前拥有的输入是两个 pandas 数据帧,我想取回一个作为输出。

通用命令行实现:

/usr/local/bin/my_command inputfileA.csv inputfileB.csv -o outputfile

可以预见的是,我当前的实现不起作用。我不知道数据帧是如何/何时通过命名管道发送到命令进程的,我将不胜感激!

import os
import StringIO
import subprocess
import pandas as pd
dfA = pd.DataFrame([[1,2,3],[3,4,5]], columns=["A","B","C"])
dfB = pd.DataFrame([[5,6,7],[6,7,8]], columns=["A","B","C"]) 

# make two FIFOs to host the dataframes
fnA = 'inputA'; os.mkfifo(fnA); ffA = open(fnA,"w")
fnB = 'inputB'; os.mkfifo(fnB); ffB = open(fnB,"w")

# don't know if I need to make two subprocesses to pipe inputs 
ppA  = subprocess.Popen("echo", 
                    stdin =subprocess.PIPE,
                    stdout=subprocess.PIPE,
                    stderr=subprocess.PIPE)
ppB  = subprocess.Popen("echo", 
                    stdin = suprocess.PIPE,
                    stdout=subprocess.PIPE,
                    stderr=subprocess.PIPE)

ppA.communicate(input = dfA.to_csv(header=False,index=False,sep="\t"))
ppB.communicate(input = dfB.to_csv(header=False,index=False,sep="\t"))


pope = subprocess.Popen(["/usr/local/bin/my_command",
                        fnA,fnB,"stdout"],
                        stdout=subprocess.PIPE,
                        stderr=subprocess.PIPE)
(out,err) = pope.communicate()

try:
    out = pd.read_csv(StringIO.StringIO(out), header=None,sep="\t")
except ValueError: # fail
    out = ""
    print("\n###command failed###\n")

os.unlink(fnA); os.remove(fnA)
os.unlink(fnB); os.remove(fnB)

【问题讨论】:

    标签: python pandas subprocess


    【解决方案1】:

    您不需要额外的进程将数据传递给子进程而不会将其写入磁盘:

    #!/usr/bin/env python
    import os
    import shutil
    import subprocess
    import tempfile
    import threading
    from contextlib import contextmanager    
    import pandas as pd
    
    @contextmanager
    def named_pipes(count):
        dirname = tempfile.mkdtemp()
        try:
            paths = []
            for i in range(count):
                paths.append(os.path.join(dirname, 'named_pipe' + str(i)))
                os.mkfifo(paths[-1])
            yield paths
        finally:
            shutil.rmtree(dirname)
    
    def write_command_input(df, path):
        df.to_csv(path, header=False,index=False, sep="\t")
    
    dfA = pd.DataFrame([[1,2,3],[3,4,5]], columns=["A","B","C"])
    dfB = pd.DataFrame([[5,6,7],[6,7,8]], columns=["A","B","C"])
    
    with named_pipes(2) as paths:
        p = subprocess.Popen(["cat"] + paths, stdout=subprocess.PIPE)
        with p.stdout:
            for df, path in zip([dfA, dfB], paths):
                t = threading.Thread(target=write_command_input, args=[df, path]) 
                t.daemon = True
                t.start()
            result = pd.read_csv(p.stdout, header=None, sep="\t")
    p.wait()
    

    cat 用于演示。您应该改用您的命令 ("/usr/local/bin/my_command")。我假设您不能使用标准输入传递数据,而必须通过文件传递输入。结果从子进程的标准输出中读取。

    【讨论】:

    • 谢谢 JF。实现很清楚,我不会想到使用线程模块来传递数据而不是子进程。我不确定我是否理解您的评论,即我不需要额外的进程来传递数据 - 这个解决方案本质上不是通过线程模块创建子进程作为子进程的替代方案吗?
    • @blackgore: 1. threading 不创建新进程 2. threading 在这里用于执行 async.io 但不是必须的;你可以使用一个选择循环来代替
    【解决方案2】:

    所以有几件事可能会让你搞砸。上一篇文章中的重要一点是像对待普通文件一样考虑这些 FIFO。除了发生的正常情况是,如果您尝试在一个进程中从管道中读取数据而不连接另一个进程在另一端写入它(反之亦然),它们会阻塞。这就是我可能会处理这种情况的方式,我会尽力描述我的想法。


    首先,当您在主进程中并尝试调用ffA = open(fnA, 'w') 时,您遇到了我上面谈到的问题——管道的另一端还没有人从中读取数据,所以发出命令后,主进程只是要阻塞。为了解决这个问题,您可能需要更改代码以删除 open() 调用:

    # make two FIFOs to host the dataframes
    fnA = './inputA';
    os.mkfifo(fnA);
    fnB = './inputB';
    os.mkfifo(fnB);
    

    好的,所以我们已经制作好管道“inputA”和“inputB”,并准备好打开读/写。为了防止像上面那样发生阻塞,我们需要启动几个子进程来调用open()。由于我对子进程库不是特别熟悉,因此我将只分叉几个子进程。

    for x in xrange(2):
    
        pid = os.fork()
        if pid == 0:
                if x == 0:
                        dfA.to_csv(open(fnA, 'w'), header=False, index=False, sep='\t')
                else:
                        dfB.to_csv(open(fnB, 'w'), header=False, index=False, sep='\t')
                exit()
        else:
                continue
    

    好的,现在我们将让这两个子进程在等待写入各自的 FIFO 时阻塞。现在我们可以运行命令连接到管道的另一端并开始读取。

    pope = subprocess.Popen(["./my_cmd.sh",
                            fnA,fnB],
                            stdout=subprocess.PIPE,
                            stderr=subprocess.PIPE)
    (out,err) = pope.communicate()
    
    try:
        out = pd.read_csv(StringIO.StringIO(out), header=None,sep="\t")
    except ValueError: # fail
        out = ""
        print("\n###command failed###\n")
    

    我发现的最后一个注释是取消链接管道似乎会删除它,因此无需调用remove()。

    os.unlink(fnA); 
    os.unlink(fnB);
    print "out: ", out
    

    在我的机器上打印语句产生:

    out:     0  1  2
    0  1  2  3
    1  3  4  5
    2  5  6  7
    3  6  7  8
    

    顺便说一句,我的命令只是一些 cat 语句:

    #!/bin/bash
    
    cat $1
    cat $2
    

    【讨论】:

    • 谢谢 NBartley。我使用直接的“cat”命令代替您包含的 shell 脚本重新创建了它。我对 os.fork 不够熟悉,无法理解 exit() 函数的存在——当我实现它时,它每次都会引发一个 SystemExit 。这是预期的行为,还是我的错字?
    • 嗯。代码是否可以正常工作?如果不是,那么我想是一个错字。
    • 但这可能是一个意想不到的系统差异——我在 Ubuntu 机器上使用 Python 2.7.3。您可以尝试的另一种方法是将 exit() 替换为 os._exit(os.EX_OK)。
    • mine 是一个类似的设置(ubuntu 14.10,python 2.7.8),所以我认为不太可能出现系统差异,但这正是我最终用作修复程序(os ._exit(1) ),它工作得很好,没有提出任何代码。
    • 再次感谢您的帮助。我已经实现并测试了这两个答案,尽可能使用共享代码,并且可以确认两者都按预期工作。我选择 J.F. Sebastian 的实现作为选择的答案,主要是因为在我的系统上使用“猫”示例进行测试时,它的效率略高(快 1.2 倍)。也就是说,应该指出 NBartley 的答案在精神和语法上可能更接近我发布的代码,因此任何低效率都可能是 OP 自己的!
    猜你喜欢
    • 2017-02-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-04-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多