【问题标题】:Condensing output with python and multiprocessing使用 python 和多处理压缩输出
【发布时间】:2011-10-25 13:51:02
【问题描述】:

所以一段时间以来,我一直在用 numpy 和多处理编写数字内容。它工作正常,但我无法收集结果。我已经通过以下方式完成了它,我将一个队列用于输入,一个用于输出。程序从输入队列中读取参数,对其进行处理,然后将结果放入输出队列。稍后在主进程中,我从队列中读出它并腌制它。像这样的:

def fun(inp,outp):
    while True:
        try:
            params = inp.get(block=False)
            results = runprocess(params)
            out.put(results,block=False)
        except Empty:
            break

稍后在主循环中我执行以下操作:

for p in processes:
    p.start()
for p in processes:
    p.join()

while True:
     try:
          out = outp.get(block=False)
          a[i] = [out]
     except Empty:
          break

 fi = open(filename,"w")
 cPickle.dump(fi,a)
 fi.close()

但不知何故,总是会发生以下两种情况之一:要么泡菜空了,要么进程挂起并保持运行,使用 0% 的 cpu(一开始它们上升到 100%,这基本上是数字运算)。对我做错了什么有什么想法吗?

好的,所以我用 Pool.map() 重做了它。让每个人都知道我是如何让它在这里工作的就是 sn-p:

    ncpus = mp.cpu_count()
    out = dict()

    params = [(a,p) for p in np.arange(0.0,2.0,0.1) for a in np.arange(0.001,2.0,0.1)]

    pool = mp.Pool(processes=ncpus)
    results = pool.map(runm,params)

    for i in results:
            sigs = np.zeros((order,order))
            sigsmf = np.zeros((order,order))
            sigseq = np.zeros((order,order))
            xs = np.array([])
            freqs = np.array([])
            [(a,p),sigs[:,:],sigsmf[:,:],sigseq[:,:],xs,freqs] = i
            out[(a,p)] = [sigs[:,:],sigsmf[:,:],sigseq[:,:],xs,freqs]
            print a, p, sigs[0,0]

像魅力一样工作,更容易实现!

感谢费迪南德!我不知道怎么做,但我想我们现在可以结束这个问题了!

【问题讨论】:

  • “i”是否在其他地方增加?
  • 您可能需要考虑使用更简单的Pool.map() 界面来为您处理所有无聊的细节:results = pool.map(runprocess, input_parameters)。见docs.python.org/library/…
  • 我实际上并没有使用 i,我使用的是字典,所以这不是问题 ;)。我将查看池界面!谢谢!
  • 这意味着我必须将 input_parameters 作为一个列表,对吧?

标签: python numpy multiprocessing multicore pickle


【解决方案1】:

您需要在至少get 调用中添加timeout,并删除block。在您当前的配置中,如果在调用get 时没有可用的项目,您将得到Empty 异常,跳出循环。如果您依赖不同的线程来填充该队列并且它没有及时填充它,它将过早退出循环并产生空结果。同样,put 可能会因为队列已满而挂起,从而挂起您的程序。

所以,使用这样的东西:

params = inp.get(timeout=1)
out.put(timeout=1)

【讨论】:

  • 其实inp队列是在调用进程之前预填充的,所以应该问题不大。
  • @AlexS:好的,那么我需要更多的上下文,比如你的线程如何运行 w.r.t。彼此。您当前的startjoin 紧接在他们之后,然后读出队列。此外,正如 Nathan 提到的他的评论,您目前将结果分配到 a[i],而不增加 i
  • 好的,我应该添加更多上下文,我已经重写它以使用 Pool.map(),因为我实际上只是分配参数空间!不过感谢您的回答!
【解决方案2】:

那是因为你有 block=False。当收集器尝试从队列中获取数据时,它不会立即在那里找到它。因此引发了 Empty 异常并跳出循环

从输入列表中获取数据时,您可以指定 block=False 作为它的预填充列表,我假设。但是,输出队列是在运行时构建的。因此,当您尝试从中获取数据时,它可能是空的,因为输入过程需要更长的时间来处理。

如果您知道输入队列的长度,那么您可以尝试无限期地阻塞输出队列 qet。如果没有,那么我建议你阻止超时。

【讨论】:

  • 实际上,我在调用 start 和 join 之间,所以至少我认为应该不成问题。我将在上面进行编辑以澄清!
  • @Alex - 这仍然是个问题。尝试根据上面的答案更改您的代码。希望它会奏效。
猜你喜欢
  • 2016-03-17
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-08-14
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多