代码之家  ›  专栏  ›  技术社区  ›  Aleph Aleph

同时监视子进程的stdout和stderr

  •  6
  • Aleph Aleph  · 技术社区  · 8 年前

    如何同时观察长时间运行的子流程的标准输出和标准错误,并在子流程生成每一行时立即对其进行处理?

    我不介意使用python3.6的异步工具在这两个流中的每一个流上创建非阻塞异步循环,但这似乎并不能解决问题。以下代码:

    import asyncio
    from asyncio.subprocess import PIPE
    from datetime import datetime
    
    
    async def run(cmd):
        p = await asyncio.create_subprocess_shell(cmd, stdout=PIPE, stderr=PIPE)
        async for f in p.stdout:
            print(datetime.now(), f.decode().strip())
        async for f in p.stderr:
            print(datetime.now(), "E:", f.decode().strip())
    
    if __name__ == '__main__':
        loop = asyncio.get_event_loop()
        loop.run_until_complete(run('''
             echo "Out 1";
             sleep 1;
             echo "Err 1" >&2;
             sleep 1;
             echo "Out 2"
        '''))
        loop.close()
    

    输出:

    2018-06-18 00:06:35.766948 Out 1
    2018-06-18 00:06:37.770187 Out 2
    2018-06-18 00:06:37.770882 E: Err 1
    

    虽然我希望它输出如下内容:

    2018-06-18 00:06:35.766948 Out 1
    2018-06-18 00:06:36.770882 E: Err 1
    2018-06-18 00:06:37.770187 Out 2
    
    2 回复  |  直到 8 年前
        1
  •  7
  •   user4815162342    8 年前

    要实现这一点,您需要一个函数,该函数将采用两个异步序列和 合并 当它们可用时,它们从一个或另一个产生结果。有了这样的功能, run 可能是这样的:

    async def run(cmd):
        p = await asyncio.create_subprocess_shell(cmd, stdout=PIPE, stderr=PIPE)
        async for f in merge(p.stdout, p.stderr):
            print(datetime.now(), f.decode().strip())
    

    类似于 merge 标准库中尚不存在,但 aiostream 外部库 provides one . 您也可以使用异步生成器编写自己的 asyncio.wait() :

    async def merge(*iterables):
        iter_next = {it.__aiter__(): None for it in iterables}
        while iter_next:
            for it, it_next in iter_next.items():
                if it_next is None:
                    fut = asyncio.ensure_future(it.__anext__())
                    fut._orig_iter = it
                    iter_next[it] = fut
            done, _ = await asyncio.wait(iter_next.values(),
                                         return_when=asyncio.FIRST_COMPLETED)
            for fut in done:
                iter_next[fut._orig_iter] = None
                try:
                    ret = fut.result()
                except StopAsyncIteration:
                    del iter_next[fut._orig_iter]
                    continue
                yield ret
    

    以上 运行 仍将在一个细节上与所需的输出不同:它不会区分输出和错误行。但这可以通过用一个指示符来修饰行来实现:

    async def decorate_with(it, prefix):
        async for item in it:
            yield prefix, item
    
    async def run(cmd):
        p = await asyncio.create_subprocess_shell(cmd, stdout=PIPE, stderr=PIPE)
        async for is_out, line in merge(decorate_with(p.stdout, True),
                                        decorate_with(p.stderr, False)):
            if is_out:
                print(datetime.now(), line.decode().strip())
            else:
                print(datetime.now(), "E:", line.decode().strip())
    
        2
  •  1
  •   user4815162342    7 年前

    在我看来,这个问题实际上有一个更简单的解决方案,至少在监视代码不需要在一个协程调用中的情况下是这样的。

    你所能做的就是产生两个独立的协程,一个用于stdout,一个用于stderr。并行运行它们将为您提供所需的语义,您可以使用 gather 等待完成:

    def watch(stream, prefix=''):
        async for line in stream:
            print(datetime.now(), prefix, line.decode().strip())
    
    async def run(cmd):
        p = await asyncio.create_subprocess_shell(cmd, stdout=PIPE, stderr=PIPE)
        await asyncio.gather(watch(p.stdout), watch(p.stderr, 'E:'))