【问题标题】:Using shared list with pathos multiprocessing raises `digest sent was rejected` error将共享列表与 pathos 多处理一起使用会引发“发送的摘要被拒绝”错误
【发布时间】:2021-02-27 15:00:13
【问题描述】:

我正在尝试按照以下代码 sn-p 使用多处理来生成复杂的、不可拾取的对象:

from multiprocessing import Manager
from pathos.multiprocessing import ProcessingPool

class Facility:

    def __init__(self):
        self.blocks = Manager().list()

    def __process_blocks(self, block):
        designer = block["designer"]
        apply_terrain = block["terrain"]
        block_type = self.__block_type_to_string(block["type"])
        block = designer.generate_block(block_id=block["id"],
                                            block_type=block_type,
                                            anchor=Point(float(block["anchor_x"]), float(block["anchor_y"]),
                                                         float(block["anchor_z"])),
                                            pcu_anchor=Point(float(block["pcu_x"]), float(block["pcu_y"]), 0),
                                            corridor_width=block["corridor"],
                                            jb_height=block["jb_connect_height"],
                                            min_boxes=block["min_boxes"],
                                            apply_terrain=apply_terrain)
        self.blocks.append(block)

    def design(self, apply_terrain=False):
        designer = FacilityBuilder(string_locator=self._string_locator, string_router=self._string_router,
                                   box_router=self._box_router, sorter=self._sorter,
                                   tracker_configurator=self._tracker_configurator, config=self._config)
        blocks = [block.to_dict() for index, block in self._store.get_blocks().iterrows()]
        for block in blocks:
            block["designer"] = designer
            block["terrain"] = apply_terrain

        with ProcessingPool() as pool:
            pool.map(self.__process_blocks, blocks)

(努力用更简单的代码重现这个,所以我展示的是实际代码)

我需要更新一个可共享变量,所以我使用multiprocessing.Manager 初始化一个类级别变量,如下所示:

self.blocks = Manager().list()

这给我留下了以下错误(仅部分堆栈跟踪):

  File "C:\Users\Paul.Nel\Documents\repos\autoPV\.autopv\lib\site-packages\dill\_dill.py", line 481, in load
    obj = StockUnpickler.load(self)
  File "C:\Users\Paul.Nel\AppData\Local\Programs\Python\Python39\lib\multiprocessing\managers.py", line 933, in RebuildProxy
    return func(token, serializer, incref=incref, **kwds)
  File "C:\Users\Paul.Nel\AppData\Local\Programs\Python\Python39\lib\multiprocessing\managers.py", line 783, in __init__
    self._incref()
  File "C:\Users\Paul.Nel\AppData\Local\Programs\Python\Python39\lib\multiprocessing\managers.py", line 837, in _incref
    conn = self._Client(self._token.address, authkey=self._authkey)
  File "C:\Users\Paul.Nel\AppData\Local\Programs\Python\Python39\lib\multiprocessing\connection.py", line 513, in Client
    answer_challenge(c, authkey)
  File "C:\Users\Paul.Nel\AppData\Local\Programs\Python\Python39\lib\multiprocessing\connection.py", line 764, in answer_challe
nge
    raise AuthenticationError('digest sent was rejected')
multiprocessing.context.AuthenticationError: digest sent was rejected

作为最后的手段,我尝试使用python 的标准ThreadPool 实现来尝试规避pickle 问题,但这也不是很顺利。我已经阅读了许多类似的问题,但还没有找到解决这个特定问题的方法。问题出在dill 还是pathosmulitprocessing.Manager 的接口方式上?

编辑:所以我设法用示例代码复制它,如下所示:

import os
import math
from multiprocessing import Manager
from pathos.multiprocessing import ProcessingPool


class MyComplex:

    def __init__(self, x):
        self._z = x * x

    def me(self):
        return math.sqrt(self._z)


class Starter:

    def __init__(self):
        manager = Manager()
        self.my_list = manager.list()

    def _f(self, value):
        print(f"{value.me()} on {os.getpid()}")
        self.my_list.append(value.me)

    def start(self):
        names = [MyComplex(x) for x in range(100)]

        with ProcessingPool() as pool:
            pool.map(self._f, names)


if __name__ == '__main__':
    starter = Starter()
    starter.start()

添加self.my_list = manager.list()时出现错误。

【问题讨论】:

    标签: python multiprocessing dill pathos


    【解决方案1】:

    所以我已经解决了这个问题。如果像 mmckerns 这样的人或比我更了解多处理的其他人可以评论为什么这是一个解决方案,我仍然会很棒。

    问题似乎在于Manager().list() 是在__init__ 中声明的。以下代码可以正常工作:

    import os
    import math
    from multiprocessing import Manager
    from pathos.multiprocessing import ProcessingPool
    
    
    class MyComplex:
    
        def __init__(self, x):
            self._z = x * x
    
        def me(self):
            return math.sqrt(self._z)
    
    
    class Starter:
    
        def _f(self, value):
            print(f"{value.me()} on {os.getpid()}")
            return value.me()
    
        def start(self):
            manager = Manager()
            my_list = manager.list()
            names = [MyComplex(x) for x in range(100)]
    
            with ProcessingPool() as pool:
                my_list.append(pool.map(self._f, names))
            print(my_list)
    
    
    if __name__ == '__main__':
        starter = Starter()
        starter.start()
    

    在这里,我将list 声明为ProcessingPool 操作的本地。如果我选择的话,我可以在之后将结果分配给一个类级别的变量。

    【讨论】:

    • 嗨@Paul:后面的代码传递了一个具有更简单依赖链的对象,因此成功的机会更大。如果您想查看传递给序列化程序的内容,可以使用以下命令打开 pickle 跟踪:import dill; dill.detect.trace(True)。当对象被序列化时,它将打印依赖链。
    • 感谢@MikeMcKerns。顺便说一句,dill 做得很好。它打开了许多因pickle 的约束而关闭的门。
    猜你喜欢
    • 1970-01-01
    • 2021-08-02
    • 1970-01-01
    • 2013-12-20
    • 2014-05-21
    • 2021-06-07
    • 2021-05-14
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多