【问题标题】:Python multiprocessing with async functions具有异步函数的 Python 多处理
【发布时间】:2019-07-31 23:07:58
【问题描述】:

我搭建了一个websocket服务器,简化版如下图:

import websockets, subprocess, asyncio, json, re, os, sys
from multiprocessing import Process

def docker_command(command_words):
    return subprocess.Popen(
        ["docker"] + command_words,
        stdout=subprocess.PIPE, stderr=subprocess.STDOUT, universal_newlines=True)

async def check_submission(websocket:object, submission:dict):
    exercise=submission["exercise"]
    with docker_command(["exec", "-w", "badkan", "grade_exercise", exercise]) as proc:
        for line in proc.stdout:
            print("> " + line)
            await websocket.send(line)

async def run(websocket, path):
    submission_json = await websocket.recv()   # returns a string
    submission = json.loads(submission_json)   # converts the string to a python dict

    ####
    await check_submission(websocket, submission)


websocketserver = websockets.server.serve(run, '0.0.0.0', 8888, origins=None)
asyncio.get_event_loop().run_until_complete(websocketserver)
asyncio.get_event_loop().run_forever()

当一次只有一个用户时,它可以正常工作。但是,当多个用户尝试使用服务器时,服务器会依次处理它们,因此以后的用户必须等待很长时间。

我尝试将标有“####”(“await check_submission...”)的行替换为:

p = Process(target=check_submission, args=(websocket, submission,))
p.start()

但是,它不起作用 - 我收到了运行时警告:“coroutine: 'check_submission' is never awaited”,并且我没有看到任何通过 websocket 的输出。

我还尝试将这些行替换为:

loop = asyncio.get_event_loop()
loop.set_default_executor(ProcessPoolExecutor())
await loop.run_in_executor(None, check_submission, websocket, submission)

但得到一个不同的错误:“can't pickle asyncio.Future objects”。

如何构建这个多处理 websocket 服务器?

【问题讨论】:

    标签: python-3.x websocket async-await python-multiprocessing python-asyncio


    【解决方案1】:

    这是我的例子,asyncio.run() 为我工作,使用多进程启动异步函数

    class FlowConsumer(Base):
        def __init__(self):
            pass
    
        async def run(self):
            self.logger("start consumer process")
            while True:
                # get flow from queue
                flow = {}
                # call flow executor get result
                executor = FlowExecutor(flow)
                rtn = FlowResult()
                try:
                    rtn = await executor.run()
                except Exception as e:
                    self.logger("flow run except:{}".format(traceback.format_exc()))
                    rtn.status = FLOW_EXCEPT
                    rtn.msg = str(e)
                self.logger("consumer flow finish,result:{}".format(rtn.dict()))
                time.sleep(1)
    
        def process(self):
            asyncio.run(self.run())
    
    
    processes = []
    consumer_proc_count = 3
    
    # start multi consumer processes 
    for _ in range(consumer_proc_count):
        # old version
        # p = Process(target=FlowConsumer().run)
        p = Process(target=FlowConsumer().process)
        p.start()
        processes.append(p)
    
    for p in processes:
        p.join()
    

    【讨论】:

    • 这可能不是解决 OP 的 问题 的最简单解决方案,但它是问题的预期答案! +1 因为我认为大多数人最终都会来到这里,因为他们遇到了 OP 声称的问题(问题所说的),而不是他确实遇到的问题。
    【解决方案2】:

    问题是subprocess.Popen不是异步的,所以check_submission在等待下一行docker输出时阻塞了事件循环。

    您根本不需要使用多处理;由于您在等待子进程时被阻塞,您只需从subprocess 切换到asyncio.subprocess

    async def docker_command(command_words):
        return await asyncio.subprocess.create_subprocess_exec(
            *["docker"] + command_words,
            stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.STDOUT)
    
    async def check_submission(websocket:object, submission:dict):
        exercise = submission["exercise"]
        proc = await docker_command(["exec", "-w", "badkan", "grade_exercise", exercise])
        async for line in proc.stdout:
            print(b"> " + line)
            await websocket.send(line)
        await proc.wait()
    

    【讨论】:

    • 在 "async for line in proc.stdout:" 行中出现错误 "AttributeError: 'generator' object has no attribute 'stdout'"。
    • @ErelSegal-Halevi 对,create_subprocess_exec 需要等待。我现在已经解决了这个问题(以及一些小问题)。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2014-03-24
    • 2013-09-06
    • 2019-11-10
    • 2021-11-19
    • 2018-12-12
    • 2020-04-08
    • 1970-01-01
    相关资源
    最近更新 更多