【问题标题】:Multiple pipes in subprocess子进程中的多个管道
【发布时间】:2015-05-04 14:55:40
【问题描述】:

我正在尝试在 ruffus 管道中使用 Sailfish,它将多个 fastq 文件作为参数。我使用python中的子进程模块执行Sailfish,但是即使我设置了shell=True,子进程调用中的<()也不起作用。

这是我要使用 python 执行的命令:

sailfish quant [options] -1 <(cat sample1a.fastq sample1b.fastq) -2 <(cat sample2a.fastq sample2b.fastq) -o [output_file]

或(最好):

sailfish quant [options] -1 <(gunzip sample1a.fastq.gz sample1b.fastq.gz) -2 <(gunzip sample2a.fastq.gz sample2b.fastq.gz) -o [output_file]

概括:

someprogram <(someprocess) <(someprocess)

我将如何在 python 中执行此操作?子流程是正确的方法吗?

【问题讨论】:

标签: python pipe subprocess named-pipes


【解决方案1】:

模拟bash process substitution:

#!/usr/bin/env python
from subprocess import check_call

check_call('someprogram <(someprocess) <(anotherprocess)',
           shell=True, executable='/bin/bash')

在 Python 中,您可以使用命名管道:

#!/usr/bin/env python
from subprocess import Popen

with named_pipes(n=2) as paths:
    someprogram = Popen(['someprogram'] + paths)
    processes = []
    for path, command in zip(paths, ['someprocess', 'anotherprocess']):
        with open(path, 'wb', 0) as pipe:
            processes.append(Popen(command, stdout=pipe, close_fds=True))
    for p in [someprogram] + processes:
        p.wait()

named_pipes(n) 在哪里:

import os
import shutil
import tempfile
from contextlib import contextmanager

@contextmanager
def named_pipes(n=1):
    dirname = tempfile.mkdtemp()
    try:
        paths = [os.path.join(dirname, 'named_pipe' + str(i)) for i in range(n)]
        for path in paths:
            os.mkfifo(path)
        yield paths
    finally:
        shutil.rmtree(dirname)

实现 bash 进程替换的另一种更可取的方式(无需在磁盘上创建命名条目)是使用 /dev/fd/N 文件名(如果可用)作为 suggested by @Dunes。在 FreeBSD 上,fdescfs(5) (/dev/fd/#) creates entries for all file descriptors opened by the process。要测试可用性,请运行:

$ test -r /dev/fd/3 3</dev/null && echo /dev/fd is available

如果失败;尝试将 /dev/fd 符号链接到 proc(5),就像在某些 Linux 上所做的那样:

$ ln -s /proc/self/fd /dev/fd

这是基于/dev/fdsomeprogram &lt;(someprocess) &lt;(anotherprocess) bash 命令的实现:

#!/usr/bin/env python3
from contextlib import ExitStack
from subprocess import CalledProcessError, Popen, PIPE

def kill(process):
    if process.poll() is None: # still running
        process.kill()

with ExitStack() as stack: # for proper cleanup
    processes = []
    for command in [['someprocess'], ['anotherprocess']]:  # start child processes
        processes.append(stack.enter_context(Popen(command, stdout=PIPE)))
        stack.callback(kill, processes[-1]) # kill on someprogram exit

    fds = [p.stdout.fileno() for p in processes]
    someprogram = stack.enter_context(
        Popen(['someprogram'] + ['/dev/fd/%d' % fd for fd in fds], pass_fds=fds))
    for p in processes: # close pipes in the parent
        p.stdout.close()
# exit stack: wait for processes
if someprogram.returncode != 0: # errors shouldn't go unnoticed
   raise CalledProcessError(someprogram.returncode, someprogram.args)

注意:在我的 Ubuntu 机器上,subprocess 代码仅适用于 Python 3.4+,尽管 pass_fds 自 Python 3.2 起可用。

【讨论】:

  • 感谢 J.F. 塞巴斯蒂安!它实际上与我之前缺少的简单子进程参数executable='/bin/bash' 一起工作。现在可以使用此调用:check_call('sailfish quant [options] &lt;(gunzip -c file1 file2) &lt;(gunzip -c file3 file4)', shell=True, executable='/bin/bash')。非常感谢你的帮助!你的回答真的超越了你——你不仅帮助我解决了我的问题,还帮助我更好地理解了 python 中的管道。
【解决方案2】:

虽然 J.F. Sebastian 提供了使用命名管道的答案,但可以使用匿名管道来做到这一点。

import shlex
from subprocess import Popen, PIPE

inputcmd0 = "zcat hello.gz" # gzipped file containing "hello"
inputcmd1 = "zcat world.gz" # gzipped file containing "world"

def get_filename(file_):
    return "/dev/fd/{}".format(file_.fileno())

def get_stdout_fds(*processes):
    return tuple(p.stdout.fileno() for p in processes)

# setup producer processes
inputproc0 = Popen(shlex.split(inputcmd0), stdout=PIPE)
inputproc1 = Popen(shlex.split(inputcmd1), stdout=PIPE)

# setup consumer process
# pass input processes pipes by "filename" eg. /dev/fd/5
cmd = "cat {file0} {file1}".format(file0=get_filename(inputproc0.stdout), 
    file1=get_filename(inputproc1.stdout))
print("command is:", cmd)
# pass_fds argument tells Popen to let the child process inherit the pipe's fds
someprogram = Popen(shlex.split(cmd), stdout=PIPE, 
    pass_fds=get_stdout_fds(inputproc0, inputproc1))

output, error = someprogram.communicate()

for p in [inputproc0, inputproc1, someprogram]:
    p.wait()

assert output == b"hello\nworld\n"

【讨论】:

  • 你的代码是:inputcmd | someproc——它不同于someproc &lt;(inputcmd)。顺便说一句,您应该调用inputproc.communicate() 而不是inputproc.wait() 来关闭父级中的inputproc.stdout.close(),这样如果someproc 过早退出,inputproc 就不会挂起。目前尚不清楚您要使用 StreamConnector 实现什么目标,但它看起来很臃肿。
  • 我的错误。我以为&lt;(cmdlist) 将一系列命令stdout 连接到消费者进程的stdin。该类旨在为流而不是文件提供类似 cat 的实用程序。答案现在简单多了。
  • /dev/fd/# 或命名管道(如果前者不可用)正是 bash 实现进程替换的方式。您应该关闭父级中的管道,以便如果inputproc1inputproc2 过早死亡; someprogram 可以早点退出。否则该解决方案应该适用于 Python 3.4+。我已将代码的异常安全版本添加到 my answer(作为练习)。
  • 关于管道我的意思是相反的:如果someprogram 过早死亡,那么在inputproc[01] 填充他们的操作系统标准输出管道缓冲区(在我的机器上约为65K)之后,父Python 脚本将挂在p.wait() ) 如果您不关闭父级中的管道 (inputproc[01].stdout)。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-10-13
  • 2015-04-27
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-04-06
相关资源
最近更新 更多