【问题标题】:Multiprocessing pool 'apply_async' only seems to call function once多处理池“apply_async”似乎只调用一次函数
【发布时间】:2015-02-09 22:21:59
【问题描述】:

我一直在关注文档以尝试了解多处理池。我想出了这个:

import time
from multiprocessing import Pool

def f(a):
    print 'f(' + str(a) + ')'
    return True

t = time.time()
pool = Pool(processes=10)
result = pool.apply_async(f, (1,))
print result.get()
pool.close()
print ' [i] Time elapsed ' + str(time.time() - t)

我正在尝试使用 10 个进程来评估函数 f(a)。我已在f 中添加了打印声明。

这是我得到的输出:

$ python pooltest.py 
f(1)
True
 [i] Time elapsed 0.0270888805389

在我看来,函数f 只被评估一次。

我可能没有使用正确的方法,但我正在寻找的最终结果是同时使用 10 个进程运行 f,并获得每个进程返回的结果。所以我会列出 10 个结果(可能相同也可能不同)。

关于多处理的文档非常混乱,要弄清楚我应该采用哪种方法并非易事,在我看来,f 在我上面提供的示例中应该运行 10 次。

【问题讨论】:

  • apply_async 并不意味着启动 多个 进程;它只是为了在池的一个进程中调用带有参数的函数。如果您希望函数被调用 10 次,则需要进行 10 次调用。
  • @JoshuaTaylor 明白了,感谢您的澄清。最好的方法是 for 循环 10 次调用,还是在 multiprocessing 模块中有更合适的工具来实现我想要实现的目标?
  • 你想每次都用相同的参数调用函数吗?
  • @JoshuaTaylor 是的,它是相同参数的 10 倍。该函数依赖于从网络中获取的其他变量,因此每次结果可能不同。
  • 啊,所以你每次都有不同的功能?然后你基本上有一个函数列表,并且你多次调用“应用这个函数”。好的

标签: python multithreading multiprocessing threadpool


【解决方案1】:

apply_async 并不意味着启动多个进程;它只是为了在池的一个进程中调用带有参数的函数。如果您希望函数被调用 10 次,则需要进行 10 次调用。

首先,请注意apply() 上的文档(已添加重点):

apply(func[, args[, kwds]])

使用参数 args 和关键字参数 kwds 调用 func。它阻塞 直到结果准备好。鉴于此块, apply_async() 更好 适合并行执行工作。 另外,func 只是 在池中的一名工人中执行。

现在,在apply_async() 的文档中:

apply_async(func[, args[, kwds[, callback[, error_callback]]]])

apply() 方法的变体,它返回一个结果对象。

两者的区别只是 apply_async 立即返回。您可以使用map() 多次调用一个函数,但如果您使用相同的输入进行调用,那么创建 same 参数的列表只是为了获得一系列长度合适。

但是,如果您使用 same 输入调用不同的函数,那么您实际上只是在调用更高阶的函数,您可以使用 mapmap_async() 来执行此操作,例如这个:

multiprocessing.map(lambda f: f(1), functions)

除了 lambda 函数不可腌制,因此您需要使用已定义的函数(请参阅 How to let Pool.map take a lambda function)。您实际上可以使用内置的apply()(不是多处理的)(尽管它已被弃用):

multiprocessing.map(apply,[(f,1) for f in functions])

自己编写也很容易:

def apply_(f,*args,**kwargs):
  return f(*args,**kwargs)

multiprocessing.map(apply_,[(f,1) for f in functions])

【讨论】:

  • 小心,我不认为 lambda 用作 lambda 函数不可腌制 (IIRC)。
  • 我了解apply 现在是如何工作的,并且从那时起它就是这样,但你第一次是对的。我正在尝试使用 same 参数运行 10 次 same 函数(尽管由于其他原因结果可能并不总是相同)。我想最终列出 10 个结果。那么最好的方法是使用for 循环吗?
  • @Juicy 我认为 map 是你想要的
  • 但是map 似乎只需要一个iterable,我想我应该只用 10 次相同的参数创建一个可迭代的列表。
  • @Juicy 查看我的更新;你有一个函数的列表;那是您应该传入的可迭代对象。您要映射的函数应该调用它的参数与 arglist
【解决方案2】:

每次您编写pool.apply_async(...) 时,它都会将该函数调用委托给在池中启动的进程之一。如果要在多个进程中调用该函数,则需要发出多个pool.apply_async调用。

注意,还有一个pool.map(和pool.map_async)函数,它将接受一个函数和一个可迭代的输入:

inputs = range(30)
results = pool.map(f, inputs)

这些函数会将函数应用到inputs 迭代中的每个输入。它尝试将“批次”放入池中,以便在池中的所有进程之间相当均匀地平衡负载。

【讨论】:

    【解决方案3】:

    如果您想在十个进程中运行一段代码,然后每个进程都退出,那么使用包含十个进程的Pool 可能不是正确的选择。

    相反,创建十个Processes 来运行代码:

    processes = []
    
    for _ in range(10):
        p = multiprocessing.Process(target=f, args=(1,))
        p.start()
        processes.append(p)
    
    for p in processes:
        p.join()
    

    multiprocessing.Pool 类旨在处理进程数和作业数不相关的情况。通常,进程数被选择为您拥有的 CPU 内核数,而作业数则要大得多。谢谢!

    【讨论】:

    • OP 想要一个您的示例未提供的返回结果。
    【解决方案4】:

    如果您出于任何特定原因没有致力于 Pool,我已经围绕 multiprocessing.Process 编写了一个函数,它可能会为您解决问题。它已发布here,但如果您需要,我很乐意将最新版本上传到 github。

    【讨论】:

      猜你喜欢
      • 2017-05-10
      • 2019-07-29
      • 2011-09-23
      • 2014-09-06
      • 2020-10-30
      • 2018-02-08
      • 1970-01-01
      • 1970-01-01
      • 2019-05-15
      相关资源
      最近更新 更多