【问题标题】:Pool workers do not complete all tasks池工作人员未完成所有任务
【发布时间】:2012-08-11 20:06:30
【问题描述】:

我有一个相对简单的 python 多处理脚本,它设置了一个工作池,通过自定义管理器将输出附加到 pandas dataframe。我发现当我在池上调用 close()/join() 时,并不是所有由 apply_async 提交的任务都已完成。

这是一个简化的示例,它提交了 1000 个作业,但只完成了一半,导致断言错误。我是否忽略了一些非常简单的事情或者这可能是一个错误?

from pandas import DataFrame
from multiprocessing.managers import BaseManager, Pool

class DataFrameResults:
    def __init__(self):
        self.results = DataFrame(columns=("A", "B")) 

    def get_count(self):
        return self.results["A"].count()

    def register_result(self, a, b):
        self.results = self.results.append([{"A": a, "B": b}], ignore_index=True)

class MyManager(BaseManager): pass

MyManager.register('DataFrameResults', DataFrameResults)

def f1(results, a, b):
    results.register_result(a, b)

def main():
    manager = MyManager()
    manager.start()
    results = manager.DataFrameResults()

    pool = Pool(processes=4)

    for (i) in range(0, 1000):
        pool.apply_async(f1, [results, i, i*i])
    pool.close()
    pool.join()

    print results.get_count()
    assert results.get_count() == 1000

if __name__ == "__main__":
    main()

【问题讨论】:

    标签: python pandas dataframe multiprocessing python-multiprocessing


    【解决方案1】:

    [EDIT]您看到的问题是由于以下代码:

    self.results = self.results.append(...)
    

    这不是原子的。所以在某些情况下,线程会在读取self.results(或在追加时)但在将新帧分配给self.results之前被中断->这个实例将丢失。

    正确的解决方案是等待使用结果对象获取结果,然后将它们全部追加到主线程中。

    【讨论】:

    • 实际上,它似乎没有帮助。我按照建议收集了所有池结果,但它们仍然中途终止。还有其他建议吗?
    • 你怎么知道他们“中途终止”?你数过返回了多少结果吗?
    • 断言检查输入数据帧的结果数是否为1000(回调数)。当我运行脚本时,数据框中的记录数通常会返回 495 到 510。每次都会产生不同的结果。
    • 有趣的是,如果我从 DataFrame 切换到列表,它会完美运行。所以可能是熊猫数据帧/多处理问题。虽然我看不到自定义管理器如何将 DataFrame 沙箱化并通过管道进行通信。
    • 因为register_result()中的赋值。查看我的编辑。
    猜你喜欢
    • 2012-06-19
    • 2022-01-19
    • 2020-01-04
    • 2019-07-21
    • 2015-07-30
    • 1970-01-01
    • 2018-10-17
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多