【问题标题】:How to keep track of status with multiprocessing and pool.map?如何使用多处理和 pool.map 跟踪状态?
【发布时间】:2016-01-16 13:01:05
【问题描述】:

我是第一次设置多处理模块,基本上,我打算按照以下方式做一些事情

from multiprocessing import pool
pool = Pool(processes=102)
results = pool.map(whateverFunction, myIterable)
print 1

据我了解,1 将在所有进程返回并完成结果后立即打印。我想对这些进行一些状态更新。实现它的最佳方法是什么?

我有点犹豫是否要打印whateverFunction()。特别是如果有大约 200 个值,我将打印 200 次类似“处理完成”的内容,这不是很有用。

我希望输出像

10% of myIterable done
20% of myIterable done

【问题讨论】:

    标签: python multiprocessing


    【解决方案1】:

    pool.map 阻塞,直到所有并发函数调用完成。 pool.apply_async 不会阻止。此外,您可以使用它的callback 参数 报告进度。每次foo 完成时,都会调用一次回调函数log_result。将foo返回的值传递给它。

    from __future__ import division
    import multiprocessing as mp
    import time
    
    def foo(x):
        time.sleep(0.1)
        return x
    
    def log_result(retval):
        results.append(retval)
        if len(results) % (len(data)//10) == 0:
            print('{:.0%} done'.format(len(results)/len(data)))
    
    if __name__ == '__main__':
        pool = mp.Pool()
        results = []
        data = range(200)
        for item in data:
            pool.apply_async(foo, args=[item], callback=log_result)
        pool.close()
        pool.join()
        print(results)
    

    产量

    10% done
    20% done
    30% done
    40% done
    50% done
    60% done
    70% done
    80% done
    90% done
    100% done
    [0, 1, 2, 3, ..., 197, 198, 199]
    

    上面的log_result函数修改了全局变量results和 访问全局变量data。您不能将这些变量传递给 log_result 因为pool.apply_async 中指定的回调函数是 总是只用一个参数调用,返回值foo

    但是,您可以创建一个闭包,这至少可以明确哪些变量 log_result 取决于:

    from __future__ import division
    import multiprocessing as mp
    import time
    
    def foo(x):
        time.sleep(0.1)
        return x
    
    def make_log_result(results, len_data):
        def log_result(retval):
            results.append(retval)
            if len(results) % (len_data//10) == 0:
                print('{:.0%} done'.format(len(results)/len_data))
        return log_result
    
    if __name__ == '__main__':
        pool = mp.Pool()
        results = []
        data = range(200)
        for item in data:
            pool.apply_async(foo, args=[item], callback=make_log_result(results, len(data)))
        pool.close()
        pool.join()
        print(results)
    

    【讨论】:

    • 太棒了。我看到您在log_result() 内使用函数范围之外的变量。我可以按照callback=lambda x: log_result(x, results) 的方式做一些事情来防止这种情况发生吗?
    • 是的,您可以,但lambda 函数也将访问其范围之外的变量。由于回调函数必须接受一个且仅一个参数,即foo 的返回值,因此无法将results(和data)作为局部变量提供给回调函数。但是,您可以使用闭包将 resultslen_data 变量置于父函数的非局部范围内,而不是全局变量。我已经编辑了上面的帖子以说明我的意思。
    • 哇,这是一个非常聪明的结构。
    猜你喜欢
    • 1970-01-01
    • 2014-04-30
    • 2019-01-30
    • 2011-07-23
    • 1970-01-01
    • 2013-10-02
    • 2015-04-11
    相关资源
    最近更新 更多