多进程嵌套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
相关产品推荐
相关产品推荐

