【问题标题】:Mulitprocess management in Python with aiomultiprocess使用 aiomultiprocess 在 Python 中进行多进程管理
【发布时间】:2020-10-31 06:21:45
【问题描述】:

我对 Python 中的多处理有疑问。我需要创建异步进程,该进程运行时间未定义,进程数也未定义。新请求一到达,就必须使用请求中的规范创建一个新流程。我们使用 ZeroMQ 进行消息传递。还有一个 Process 是从一开始就开始,只有在整个脚本终止时才结束。

现在我正在寻找一种解决方案,如何等待所有进程,同时能够添加其他进程。

asyncio.gather()

这是我的第一个想法,但它在被调用之前需要进程列表。

class Object:
  def __init__(self, var):
     self.var = var

  async def run(self):
      *do async things*

class object_controller:
  
  def __init__(self):
     self.ctx = zmq.Context()
     self.socket = self.ctx.socket(zmq.PULL)
     self.socket.connect("tcp://127.0.0.1:5558")

     self.static_process = AStaticProcess()
     self.sp = aiomultiprocess.Process(target=self.static_process.run)
     self.sp.start()
     #here I need a good way to await this process


  def process(self, var):
    object = Object(var)
    process = aiomultiprocess.Process(target=object.run)
    process.start()
  
  def listener(self)
    while True:
      msg = self.socket.recv_pyobj()
      # here I need to find a way how I can start and await this process while beeing able to 
      # receive additional request, which result in additional processes which need to be awaited

这是一些希望能解释我的问题的代码。我需要一种等待进程的收集器。

初始化之后,对象和控制器之间没有交互,只有zeroMQ(静态进程和变量进程之间)。也没有回报。

【问题讨论】:

    标签: python python-asyncio zeromq


    【解决方案1】:

    如果您需要在同时等待新进程的同时启动进程,而不是显式调用await 来了解进程何时完成,请使用asyncio.create_task() 让它们在后台执行。这将返回一个Task 对象,该对象有一个add_done_callback 方法,您可以在该过程完成时使用它来做一些工作:

    class Object:
      def __init__(self, var):
         self.var = var
    
      async def run(self):
          *do async things*
    
    class object_controller:
      
      def __init__(self):
         self.ctx = zmq.Context()
         self.socket = self.ctx.socket(zmq.PULL)
         self.socket.connect("tcp://127.0.0.1:5558")
    
         self.static_process = AStaticProcess()
         self.sp = aiomultiprocess.Process(target=self.static_process.run)
         self.sp.start()
         asyncio.create_task(self.sp.join() self.handle_proc_finished)
    
    
      def process(self, var):
        object = Object(var)
        process = aiomultiprocess.Process(target=object.run)
        process.start()
      
      def listener(self)
        while True:
          msg = self.socket.recv_pyobj()
          process = aiomultiprocess.Process(...)
          process.start()
          t = asyncio.create_task(process.join())
          t.add_done_callback(self.handle_other_proc_finished)
    
      def handle_proc_finished(self, task):
         # do something
    
      def handle_other_proc_finished(self, task):
        # do something else
    

    如果你想避免使用回调,你也可以传递 create_task 一个你自己定义的协程,它等待进程完成,然后做任何需要做的事情。

    self.sp.start()
    asyncio.create_task(wait_for_proc(self.sp))
    
    async def wait_for_proc(proc):
       await proc.join()
       # do other stuff
    

    【讨论】:

    • 首先感谢您的回复!这听起来就像我正在寻找的解决方案,但不幸的是我收到错误RuntimeError: no running event loop sys:1: RuntimeWarning: coroutine 'Process.join' was never awaited 我不确定如果错过了会怎样。非常欢迎任何帮助!
    • @N.Icenstein 没有完整的复制器很难说,但这可能是因为您在启动事件循环之前创建了一个任务。您可能需要稍微重构一下,以便在您在任何地方调用 add_task 之前启动事件循环(尤其是 object_controller 构造函数?)
    • 我现在正在构造函数中创建一个事件循环,我现在正在使用asyncio.ensure_future,到目前为止效果很好(我目前正在进行一些测试以确保一切都正确)跨度>
    • @N.Icenstein 我建议在调用构造函数之前启动事件循环——从设计的角度来看它看起来更简洁——但我很高兴它解决了问题。
    【解决方案2】:

    您需要为流程创建任务列表或未来对象。此外,您不能在等待其他任务时将进程添加到事件循环中

    【讨论】:

    • 所以理论上我必须为每个新任务创建一个单独的事件循环?这甚至可能和实用吗?理论上是否可以在 while 循环结束时调用一个方法来检查所有未来的对象然后继续?
    • 这会破坏异步的目的。 Python 实际上并不是多线程的。要实现异步,您必须事先知道要异步运行哪些任务并将这些任务收集到一个列表中并等待 asyncio.wait(task list)。您可以在 while 循环中创建任务列表并等待它们,但最终您不会完全异步。您可以尝试使用 Asyncio 队列。见:stackoverflow.com/questions/28115253/…
    • 您必须事先知道要异步运行哪些任务 - 这不是真的,您可以在执行期间的任何时候添加新任务。 gatherwait 只是方便的函数来等待一些事情完成;人们可以轻松(并且经常这样做)等待抽象的结束事件并根据需要生成任务。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-03-18
    • 1970-01-01
    • 2018-05-02
    • 2021-10-17
    • 1970-01-01
    相关资源
    最近更新 更多