【问题标题】:Fill a Queue with Objects from several data loaders using multiprocessing使用多处理使用来自多个数据加载器的对象填充队列
【发布时间】:2020-03-09 08:14:46
【问题描述】:

我从事机器学习输入管道的工作。我编写了一个数据加载器,它从一个大的 .hdf 文件中读取数据并返回切片,每个切片大约需要 2 秒。因此,我想使用一个队列,它从多个数据加载器中接收对象,并可以通过下一个函数(如生成器)从队列中返回单个对象。此外,填充队列的进程应该以某种方式在后台运行,当队列未满时重新填充队列。我没有让它正常工作。它与单个数据加载器一起工作,给了我 4 倍相同的切片..

import multiprocessing as mp

class Queue_Generator():
    def __init__(self, data_loader_list):
        self.pool = mp.Pool(4)
        self.data_loader_list = data_loader_list
        self.queue = mp.Queue(maxsize=16)
        self.pool.map(self.fill_queue, self.data_loader_list)
    def fill_queue(self,gen):
        self.queue.put(next(gen))
    def __next__(self):
        yield self.queue.get()

我从中得到的: NotImplementedError:池对象不能在进程之间传递或腌制 提前致谢

【问题讨论】:

    标签: python multiprocessing generator


    【解决方案1】:

    不太清楚你为什么在__next__ 中使用yielding,这对我来说看起来不太正确。 __next__ 应该返回一个值,而不是生成器对象。

    这是一种简单的方法,您可以将并行函数的结果作为生成器返回。它可能满足也可能不满足您的特定要求,但可以进行调整以适应。它将继续处理 data_loader_list 直到用尽。例如,与始终在 Queue 中保留 4 个项目相比,这可能会占用大量内存。

    import multiprocessing as mp
    
    
    def read_lines(data_loader):
        from time import sleep
        sleep(2)
        return f'did something with {data_loader}'
    
    
    def make_gen(data_loader_list):
        with mp.Pool(4) as pool:
            for result in pool.imap(read_lines, data_loader_list):
                yield result
    
    
    if __name__ == '__main__':
        data_loader_list = [i for i in range(15)]
        result_generator = make_gen(data_loader_list)
        print(type(result_generator))
    
        for i in result_generator:
            print(i)
    

    使用imap 意味着可以在生成结果时对其进行处理。 mapmap_async 将阻塞在 for 循环中,直到所有结果都准备好。请参阅this question 了解更多信息。

    【讨论】:

      【解决方案2】:

      您的具体错误意味着当您将类方法传递给池时,您不能将池作为类的一部分。我的建议可能如下:

      import multiprocessing as mp
      from queue import Empty
      
      
      class QueueGenerator(object):
          def __init__(self, data_loader_list):
              self.data_loader_list = data_loader_list
              self.queue = mp.Queue(maxsize=16)
      
          def __iter__(self):
              processes = list()
              for _ in range(4):
                  pr = mp.Process(target=fill_queue, args=(self.queue, self.data_loader_list))
                  pr.start()
                  processes.append(pr)
              return self
      
          def __next__(self):
              try:
                  return self.queue.get(timeout=1) # this should have a value, otherwise your loop will never stop. make it something that ensures your processes have enough time to update the queue but not too long that your program freezes for an extended period of time after all information is processed
              except Empty:
                  raise StopIteration
      
      # have fill queue as a separate function
      def fill_queue(queue, gen):
          while True:
              try:
                  value = next(gen)
                  queue.put(value)
              except StopIteration: # assumes the given data_loader_list is an iterator
                  break
          print('stopping')
      
      
      gen = iter(range(70))
      
      qg = QueueGenerator(gen)
      
      
      for val in qg:
          print(val)
      # test if it works several times:
      for val in qg:
          print(val)
      

      我认为您要解决的下一个问题是让 data_loader_list 成为在每个单独的进程中提供新信息的东西。但是由于您没有提供任何有关此的信息,因此我无法为您提供帮助。然而,上面确实为您提供了一种让进程填充您的队列的方法,然后将其作为迭代器传递出去。

      【讨论】:

      • 看起来不错。我今天试试,非常感谢。 data_loader 对象返回 3D 图像的随机裁剪,因此它们应该在每个过程中提供新信息。
      • 像魅力一样工作。感谢您的快速解决方案。我遇到的下一个问题是进程返回了相同的图像裁剪,因为 np.random.randint() 基于时间生成伪随机数并且进程同步运行。因此,必须为每个 rng 设置不同的种子,如 stackoverflow.com/questions/24345637/… 中所述
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2016-11-15
      • 1970-01-01
      • 1970-01-01
      • 2019-10-02
      • 1970-01-01
      • 2017-09-11
      • 1970-01-01
      相关资源
      最近更新 更多