【问题标题】:Can asyncio subprocess be used with contextmanager?asyncio 子进程可以与 contextmanager 一起使用吗?
【发布时间】:2019-08-01 16:06:02
【问题描述】:

在 python (3.7+) 中,我尝试将子进程作为上下文管理器运行,同时异步流式传输潜在的大量标准输出。问题是我似乎无法让 contextmanager 的主体与 stdout 回调异步运行。我曾尝试使用线程,在那里运行异步函数,但后来我无法弄清楚如何将 Process 对象返回到 contextmanager。

所以问题是:如何在主线程运行时从主线程中的上下文管理器产生异步进程对象?也就是说,我想在 open_subprocess() 在下面的代码中完成运行之前从它产生已经和当前正在运行的进程。

import asyncio
import contextlib

async def read_stream(proc, stream, callback):
    while proc.returncode is None:
        data = await stream.readline()
        if data:
            callback(data.decode().rstrip())
        else:
            break

async def stream_subprocess(cmd, *args, stdout_callback=print):
    proc = await asyncio.create_subprocess_exec(
        cmd,
        *args,
        stdout=asyncio.subprocess.PIPE)
    read = read_stream(proc, proc.stdout, stdout_callback)
    await asyncio.wait([read])
    return proc

@contextlib.contextmanager
def open_subprocess(cmd, *args, stdout_callback=print):
    proc_coroutine = stream_subprocess(
        cmd,
        *args,
        stdout_callback=stdout_callback)
    # The following blocks until proc has finished
    # I would like to yield proc while it is running
    proc = asyncio.run(proc_coroutine)
    yield proc
    proc.terminate()

if __name__ == '__main__':
    import time

    def stdout_callback(data):
        print('STDOUT:', data)

    with open_subprocess('ping', '-c', '4', 'localhost',
                         stdout_callback=stdout_callback) as proc:
        # The following code only runs after proc completes
        # but I would expect these print statements to
        # be interleaved with the output from the subprocess
        for i in range(2):
            print(f'RUNNING SUBPROCESS {proc.pid}...')
            time.sleep(1)

    print(f'RETURN CODE: {proc.returncode}')

【问题讨论】:

  • 你需要 asynccontextmanager docs.python.org/3/library/…
  • @altunyurt 这看起来确实很有希望,但是等待对asyncio.create_subprocess_exec() 的调用会阻塞并阻止进入 asynccontextmanager。

标签: python python-3.x subprocess python-asyncio contextmanager


【解决方案1】:

Asyncio 通过挂起任何看起来可能会阻塞的东西来提供并行执行。为此,所有代码都必须在回调或coroutines 内,并避免调用像time.sleep() 这样的阻塞函数。除此之外,您的代码还有一些其他问题,例如 await asyncio.wait([x]) 等同于 await x,这意味着在所有流读取完成之前open_subprocess 不会产生。

构建代码的正确方法是将顶级代码移动到async def 并使用异步上下文管理器。例如:

import asyncio
import contextlib

async def read_stream(proc, stream, callback):
    while proc.returncode is None:
        data = await stream.readline()
        if data:
            callback(data.decode().rstrip())
        else:
            break

@contextlib.asynccontextmanager
async def open_subprocess(cmd, *args, stdout_callback=print):
    proc = await asyncio.create_subprocess_exec(
        cmd, *args, stdout=asyncio.subprocess.PIPE)
    asyncio.create_task(read_stream(proc, proc.stdout, stdout_callback))
    yield proc
    if proc.returncode is None:
        proc.terminate()
        await proc.wait()

async def main():
    def stdout_callback(data):
        print('STDOUT:', data)

    async with open_subprocess('ping', '-c', '4', 'localhost',
                               stdout_callback=stdout_callback) as proc:
        for i in range(2):
            print(f'RUNNING SUBPROCESS {proc.pid}...')
            await asyncio.sleep(1)

    print(f'RETURN CODE: {proc.returncode}')

asyncio.run(main())

如果您坚持混合使用同步和异步代码,则需要通过在单独的线程中运行 asyncio 事件循环来完全分离它们。那么你的主线程将无法直接访问像proc 这样的异步对象,因为它们不是线程安全的。您需要始终使用call_soon_threadsaferun_coroutine_threadsafe 与事件循环进行通信。

这种方法很复杂,需要线程间通信和摆弄事件循环,所以除了作为学习练习外,我不建议这样做。更不用说如果你使用另一个线程,你根本不需要 asyncio ——你可以直接在另一个线程中发出同步调用。但话虽如此,这里有一个可能的实现:

import asyncio
import contextlib
import concurrent.futures
import threading

async def read_stream(proc, stream, callback):
    while proc.returncode is None:
        data = await stream.readline()
        if data:
            callback(data.decode().rstrip())
        else:
            break

async def stream_subprocess(cmd, *args, proc_data_future, stdout_callback=print):
    try:
        proc = await asyncio.create_subprocess_exec(
            cmd, *args, stdout=asyncio.subprocess.PIPE)
    except Exception as e:
        proc_data_future.set_exception(e)
        raise
    proc_data_future.set_result({'proc': proc, 'pid': proc.pid})
    await read_stream(proc, proc.stdout, stdout_callback)
    return proc

@contextlib.contextmanager
def open_subprocess(cmd, *args, stdout_callback=print):
    loop = asyncio.new_event_loop()
    # needed to use asyncio.subprocess outside the main thread
    asyncio.get_child_watcher().attach_loop(loop)
    threading.Thread(target=loop.run_forever).start()
    proc_data_future = concurrent.futures.Future()
    loop.call_soon_threadsafe(
        loop.create_task,
        stream_subprocess(cmd, *args,
                          proc_data_future=proc_data_future,
                          stdout_callback=stdout_callback))
    proc_data = proc_data_future.result()
    yield proc_data
    async def terminate(proc):
        if proc.returncode is None:
            proc.terminate()
            await proc.wait()
    asyncio.run_coroutine_threadsafe(terminate(proc_data['proc']), loop).result()
    proc_data['returncode'] = proc_data['proc'].returncode
    loop.call_soon_threadsafe(loop.stop)

if __name__ == '__main__':
    import time

    def stdout_callback(data):
        print('STDOUT:', data)

    with open_subprocess('ping', '-c', '4', 'localhost',
                         stdout_callback=stdout_callback) as proc_data:
        for i in range(2):
            print(f'RUNNING SUBPROCESS {proc_data["pid"]}...')
            time.sleep(1)

    print(f'RETURN CODE: {proc_data["returncode"]}')

【讨论】:

  • 此答案提供的关键见解是无需等待已创建的任务。我认为我必须明确等待应该在事件循环中运行的所有任务——特别是从 read_stream 协程创建的任务。我仍然不确定这是为什么。
  • @Johann 你的直觉并没有完全错位——最终等待生成的任务确实是个好主意,否则任务引发的未处理异常将丢失。但确实不需要立即等待任务 - 在这种情况下,在 yield 之后等待从 read_stream() 创建的任务可能是个好主意。有关该主题的更长论文,请参阅this article
【解决方案2】:

@contextlib.asynccontextmanagerProcess.wait() 例程的方法(等待子进程终止,设置并返回returncode 属性):

import asyncio
import contextlib

async def read_stream(proc, stream, callback):
    while proc.returncode is None:
        data = await stream.readline()
        if not data:
            break
        callback(data.decode().rstrip())


async def stream_subprocess(cmd, *args, stdout_callback=print):
    proc = await asyncio.create_subprocess_exec(cmd, *args,
                                                stdout=asyncio.subprocess.PIPE)
    await read_stream(proc, proc.stdout, stdout_callback)
    return proc


@contextlib.asynccontextmanager
async def open_subprocess(cmd, *args, stdout_callback=print):
    try:
        proc = await stream_subprocess(cmd, *args, stdout_callback=stdout_callback)
        yield proc
    finally:
        await proc.wait()

if __name__ == '__main__':
    import time

    def stdout_callback(data):
        print('STDOUT:', data)


    async def main():
        async with open_subprocess('ping', '-c', '4', 'localhost',
                                   stdout_callback=stdout_callback) as proc:
            # The following code only runs after proc completes
            for i in range(2):
                print(f'RUNNING SUBPROCESS {proc.pid}...')
                time.sleep(1)

        print(f'RETURN CODE: {proc.returncode}')

    asyncio.run(main())

示例运行输出:

STDOUT: PING localhost (127.0.0.1): 56 data bytes
STDOUT: 64 bytes from 127.0.0.1: icmp_seq=0 ttl=64 time=0.048 ms
STDOUT: 64 bytes from 127.0.0.1: icmp_seq=1 ttl=64 time=0.074 ms
STDOUT: 64 bytes from 127.0.0.1: icmp_seq=2 ttl=64 time=0.061 ms
STDOUT: 64 bytes from 127.0.0.1: icmp_seq=3 ttl=64 time=0.067 ms
STDOUT: 
STDOUT: --- localhost ping statistics ---
STDOUT: 4 packets transmitted, 4 packets received, 0.0% packet loss
STDOUT: round-trip min/avg/max/stddev = 0.048/0.062/0.074/0.010 ms
RUNNING SUBPROCESS 35439...
RUNNING SUBPROCESS 35439...
RETURN CODE: 0

Process finished with exit code 0

【讨论】:

  • 我想我不清楚。我期待“RUNNING SUBPROCESS”行与“STDOUT”行交错。在这里,直到进程完成后才会输入上下文的主体。
  • @Johann,你没有提到“交错”,你没有 asynccontextmanager - 你现在有了。要获得复杂的进程间通信,您需要详细说明所有规则并呈现预期的输出
  • 我认为我很清楚这句话(代码中的注释)“我想在运行时产生 proc”。在我和您的示例代码中,上下文仅在过程完成后才会产生。请注意,我们的两个示例都给出了相同的输出。我会将该评论复制到问题的散文中。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2012-11-22
  • 1970-01-01
  • 2019-05-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2011-08-15
相关资源
最近更新 更多