You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

多进程嵌套DictProxy对象调用.items()等方法时偶发BrokenPipeError崩溃问题求助

多进程嵌套DictProxy对象调用.items()等方法时偶发BrokenPipeError崩溃问题求助

各位大佬好,我最近在做一个多进程任务跟踪功能,用multiprocessing.Manager()创建了嵌套的共享字典(外层和内层都是DictProxy对象),用来记录每个命令的执行状态和所在队列。但遇到了一个偶发的崩溃问题:调用字典的.items()、.get()方法时,程序大概每5次运行就会崩溃一次,报错信息主要是BrokenPipeError: [Errno 32] Broken pipe和File "<string>", line 2, in keys这类。

崩溃主要出现在traverseRootAndNestedKey方法里,这个方法是用来遍历嵌套字典、更新命令状态的。目前为了排查问题,我已经把嵌套字典的值简化成了单个字符,下面是完整的代码:

import multiprocessing
import multiprocessing.managers
import os, sys
import pprint

class ProgressDictionary:
    def __init__(self):
        self.dictionary = {}
        # E.g.    k[0] = {'Root Key A': Manager().dict({'Nested key A_1': [1, 0],
        #                                               'Nested key A_2': [1, 0],
        #                                               'Nested key A_3': [1, 0],
        #                                               'Nested key A_4': [1, 0]})}
        self.mprocDictionary = multiprocessing.Manager().dict()
        # Will need to be rebuilt each time self.dictionary is .updated()
        # E.g. mprocDictionary.update(<......>))
        #         = Manager().dict{'Root Key A': Manager().dict({'Nested key A_1': [1, 0],
        #                                                                   ...
        #                                                        'Nested key A_4': [1, 0]}),
        #                          'Root Key B': Manager().dict({'Nested key B_1': [1, 0],
        #                                                                   ...
        #                                                        'Nested key B_4': [1, 0]})}
        self.container = multiprocessing.Manager().dict()   # Not used

    def buildMprocDict(self, rootKey, listOfNestedKeys):
        self.dictionary.update({rootKey: self.addNest(listOfNestedKeys)})
        # Debug
        # pprint.PrettyPrinter(width=128, sort_dicts=False).pprint(self.dictionary)
        # No braces as intended   -> 'Root Key A': <DictProxy object, typeid 'dict' at 0x72335a1908e0>
        self.mprocDictionary.update(self.dictionary)
        # Debug
        # print(type(self.mprocDictionary))
        # {'Root Key A': <DictProxy object, typeid 'dict' at 0x79b399d989a0>,
        #  'Root Key B': <DictProxy object, typeid 'dict' at 0x79b399d98eb0>}
        # type(self.mprocDictionary): <class 'multiprocessing.managers.DictProxy'>

        # Not used
        self.container.update({i : multiprocessing.Manager().dict()
                                    for i in self.mprocDictionary.keys()})
        # Debug
        # for k, v in self.mprocDictionary.items():
        #     print(f'{k} ---:--- {v}')
        # print(f'K43: {[i for i in self.mprocDictionary.keys()]}')
        # ['Root Key A', 'Root Key B']
        # print('L26: class function exit')
    def addNest(self, listOfNestedKeys):
        # Encapsulate nested dicts
        return multiprocessing.Manager().dict({i: 'A' for i in listOfNestedKeys})
    def traverseNestedKeyOnly(self, rootKey, nestedKey):
        for k, v in dict(self.mprocDictionary.get(rootKey, {}).items()):
            if k == nestedKey:
                return k
    def traverseRootAndNestedKey(self, nestedKey):

        for k in self.mprocDictionary.keys():
            print(f'L56 {self.mprocDictionary.get(k).items()}')        # Works
            print(f'L57 {type(dict(self.mprocDictionary.get(k)))}')    #  Will occasionally generate error (1 in 5 chance)
                                                                    # multiprocessing.managers.DictProxy'> w/out(dict())
            print(f'L58 {(dict(self.mprocDictionary.get(k))).items()}') # File "<string>", line 2, in keys
            sys.exit()
            for k2, v in dict(self.mprocDictionary.get(k)).items():     # crashes, File "<string>", line 2, in keys
                print(f'L64 {k2} -:- {v}')
        sys.exit()
        # temp_dict = {k: dict(v) if isinstance(v, multiprocessing.managers.DictProxy) else v
        #              for k, v in self.mprocDictionary.items()}
        #pprint.PrettyPrinter(width=96, sort_dicts=False).pprint(temp_dict)
        sys.exit()
        for k, v in self.mprocDictionary.items():  #File "<string>", line 2, in __getitem__
            for k2, v2 in v.items():    # We only know value, but not key
                # In pyCharm, v.items() gives warning: Unresolved attribute reference 'items' for class 'str'
                if k2 == nestedKey:
                    return [k, k2]
        #if nested key is not found due to being deleted earlier for whatever reason
        return None
    def setMemAddress(self, nestedKey, pid, debug_helper):
        # *args used, incase root key is never supplied.
        print(f'L64 argument show: {debug_helper}')   # No show
        k, k2 = self.traverseRootAndNestedKey(nestedKey)      # for k, v in self.mprocDictionary.items():
                                                              # File "<string>", line 2, in items
        # Sometimes will cut off here, and won't show the following statements
        # Converting dictProxy to dict:  for k, v in dict(self.mprocDictionary).items(),
        #                                         ...In traverseRootAndNestedKey() method .did not resolve this
        print(f'L65 key: {k}, val: {k2}')                  # display correct values
        #print(f'L68 {self.mprocDictionary[k]}')  # <DictProxy object, typeid 'dict' at 0x796feca0fd60; '__str__()' failed>
        #print(f'L69 {dict(self.mprocDictionary[k])}')   # File "<string>", line 2, in __getitem__
        if k is not None:
            # error here, nothing is shown
            temp = self.mprocDictionary[k][k2]
            print(f'L93 {temp}')                     # Will sometimes show None, other times 'A' w/Multiproc
            temp = pid
            # Alternative method
            self.mprocDictionary[k][k2].update(temp)
    def getMemAddress(self, rootKey, nestedKey):
        k = self.traverseNestedKeyOnly(rootKey, nestedKey)
        try:
            return self.mprocDictionary[rootKey][k]
        except IndexError:
            return None
    def printDictionary(self):
        pprint.PrettyPrinter(width=96, sort_dicts=False).pprint(self.dictionary)
    def printMprocDictionary(self):
        for k, v in self.mprocDictionary.items():
            for k2, v2 in v.items():
                print(f'L78 {k2}  :  {v2}')

splitCMDqueue = multiprocessing.Queue()

qList = [multiprocessing.Queue() for i in range(2)]
listCMDs = []
listCMDs.append(['''system cmd a_1''',
                 '''system cmd a_2''',
                 '''system cmd a_3''',
                 '''system cmd a_4'''])
listCMDs.append(['''system cmd b_1''',
                 '''system cmd b_2''',
                 '''system cmd b_3''',
                 '''system cmd b_4'''])
listCMDcondition = ['''condition_system_cmd_1''',
                    '''condition_system_cmd_2''']

# {'Root Key A': {
#      'Nested key A_1': [1, 0],
#      'Nested key A_2': [1, 0],
#      'Nested key A_3': [1, 0],
#      'Nested key A_4': [1, 0]},
#  'Root Key B': {
#      'Nested key B_1': [1, 0],
#      'Nested key B_2': [1, 0],
#      'Nested key B_3': [1, 0],
#      'Nested key B_4': [1, 0]}}
def worker(sharedClass, nK):
    sharedClass.setMemAddress(nK, '0x12345FA', 'Call from Line 110')
    #print("L48") # Shows when previous statement is commented out
    print(f'L130, {sharedClass.getMemAddress(listCMDcondition[0], nK)}')

jobs = []
runtimeCommandProgress =  ProgressDictionary()

if __name__ == '__main__':

    for idx, i in enumerate(listCMDs):
        runtimeCommandProgress.buildMprocDict(listCMDcondition[idx], i)  # works/success
        for j in i:
            qList[idx].put(j)

    example_cmd = listCMDs[0][1]  # Nested key A_2, system cmd a_2
    # Both work
    # runtimeCommandProgress.setMemAddress(example_cmd, '0x12345FA', 'L136')
    # print(f'L133 {runtimeCommandProgress.getMemAddress(listCMDcondition[0], example_cmd)}')

    p = multiprocessing.Process(target=worker, args=(runtimeCommandProgress, example_cmd))

    jobs.append(p)
    for p in jobs:
        p.start()

我自己尝试过的排查方向:

  • 尝试把DictProxy转成普通字典再操作,但转换过程偶发还是会触发错误;
  • 单独调用.get()有时候没问题,但组合使用.get()和.items()就容易崩溃;
  • 报错位置不固定,有时候在.keys(),有时候在.items(),感觉是IPC通信的问题,但不知道怎么彻底解决。

有没有大佬遇到过类似的问题?或者能帮我分析下问题出在哪,怎么修复吗?


备注:内容来源于stack exchange,提问作者leipsohul

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.14 16:48:00