Python并发场景下嵌套defaultdict的管理方案探讨
多进程中嵌套defaultdict的管理问题
已有一些讨论关注多进程中非嵌套defaultdict的行为,但管理如defaultdict(list)这类嵌套结构并非易事,更不用说defaultdict(lambda: defaultdict(list))这类更复杂的嵌套结构。
尝试方案1:注册defaultdict到自定义Manager
import concurrent.futures from collections import defaultdict import multiprocessing as mp from multiprocessing.managers import BaseManager, DictProxy, ListProxy import numpy as np def called_function1(hey, i, yo): yo[i].append(hey) class EmbeddedManager(BaseManager): pass def func1(): emanager = EmbeddedManager() emanager.register('defaultdict', defaultdict, DictProxy) emanager.start() ddict = emanager.defaultdict(list) with concurrent.futures.ProcessPoolExecutor(8) as executor: for i in range(10): ind = np.random.randint(2) executor.submit(called_function1, i, ind, ddict) for k, v in ddict.items(): print(k, v) emanager.shutdown()
执行结果:
func1() 1 [] 0 []
问题:仅保留键,内部列表内容未被正确管理。
尝试方案2:手动判断键并创建普通列表
def called_function2(hey, i, yo): if i not in yo: yo[i] = [] yo[i].append(hey) def func2(): manager = mp.Manager() ddict = manager.dict() with concurrent.futures.ProcessPoolExecutor(8) as executor: for i in range(10): ind = np.random.randint(2) executor.submit(called_function2, i, ind, ddict) for k, v in ddict.items(): print(k, v)
执行结果:
func2() 1 [] 0 []
问题:普通列表无法跨进程共享,内容未被保留。
尝试方案3:提前创建受管理的列表
def called_function3(hey, i, yo): yo[i].append(hey) def func3(): manager = mp.Manager() ddict = manager.dict() with concurrent.futures.ProcessPoolExecutor(8) as executor: for i in range(10): ind = np.random.randint(2) if ind not in ddict: ddict[ind] = manager.list() executor.submit(called_function2, i, ind, ddict) for k, v in ddict.items(): print(k, v)
执行结果:
func3() 0 [0, 2, 3, 4, 6, 8] 1 [1, 5, 7, 9]
问题:需要提前预知所有可能的键,灵活性不足。
尝试方案4:传递Manager到子进程
def called_function4(hey, i, yo, man): if i not in yo: yo[i] = man.list() yo[i].append(hey) def func4(): manager = mp.Manager() ddict = manager.dict() with concurrent.futures.ProcessPoolExecutor(8) as executor: futures = [] for i in range(10): ind = np.random.randint(2) futures.append(executor.submit(called_function2, i, ind, ddict, manager)) for f in concurrent.futures.as_completed(futures): print(f.result()) for k, v in ddict.items(): print(k, v)
执行报错:
func4() TypeError: Pickling an AuthenticationString object is disallowed for security reasons
问题:Manager无法被序列化传递到子进程。
尝试方案5:在子进程内新建Manager
def called_function5(hey, i, yo): if i not in yo: yo[i] = mp.Manager().list() yo[i].append(hey) def func5(): manager = mp.Manager() ddict = manager.dict() with concurrent.futures.ProcessPoolExecutor(8) as executor: futures = [] for i in range(10): ind = np.random.randint(2) futures.append(executor.submit(called_function5, i, ind, ddict)) for f in concurrent.futures.as_completed(futures): print(f.result()) for k, v in ddict.items(): print(k, v)
执行报错:
func5() BrokenPipeError: [Errno 32] Broken pipe
问题:子进程内创建的Manager无法和主进程的共享字典正常交互。
请问是否存在更优的实现方式?
解决方案
方法1:自定义支持嵌套的受管理defaultdict
通过扩展BaseManager,注册能生成嵌套受管理结构的工厂函数,确保所有层级的容器都是进程安全的代理对象:
import concurrent.futures from collections import defaultdict import multiprocessing as mp from multiprocessing.managers import BaseManager, DictProxy import numpy as np # 定义嵌套defaultdict的工厂函数 def nested_defaultdict(): return defaultdict(mp.Manager().list) class NestedManager(BaseManager): pass # 注册工厂函数,指定返回DictProxy确保字典本身是受管理的 NestedManager.register('nested_defaultdict', nested_defaultdict, DictProxy) def called_function(hey, i, yo): yo[i].append(hey) def func(): manager = NestedManager() manager.start() ddict = manager.nested_defaultdict() with concurrent.futures.ProcessPoolExecutor(8) as executor: for i in range(10): ind = np.random.randint(2) executor.submit(called_function, i, ind, ddict) for k, v in ddict.items(): print(k, list(v)) # 转换为普通列表查看内容 manager.shutdown() if __name__ == "__main__": func()
执行后能正确保留所有进程写入的内容,且无需提前预知键。
方法2:进程安全字典+锁实现defaultdict逻辑
用进程安全的锁保证原子操作,动态创建受管理列表,无需自定义Manager:
import concurrent.futures import multiprocessing as mp import numpy as np def called_function(hey, i, shared_dict, lock): with lock: if i not in shared_dict: shared_dict[i] = mp.Manager().list() shared_dict[i].append(hey) def func(): manager = mp.Manager() shared_dict = manager.dict() lock = manager.Lock() # 进程安全的锁 with concurrent.futures.ProcessPoolExecutor(8) as executor: for i in range(10): ind = np.random.randint(2) executor.submit(called_function, i, ind, shared_dict, lock) for k, v in shared_dict.items(): print(k, list(v)) if __name__ == "__main__": func()
锁是必须的,避免多进程同时操作时出现竞争条件。
方法3:本地收集+主进程合并(高性能方案)
如果不需要实时共享数据,让子进程本地收集结果,最后在主进程合并,避免进程间通信开销:
import concurrent.futures from collections import defaultdict import numpy as np def called_function(hey, i): # 子进程本地处理,返回键和值 return i, hey def func(): local_dict = defaultdict(list) with concurrent.futures.ProcessPoolExecutor(8) as executor: results = executor.map(called_function, range(10), [np.random.randint(2) for _ in range(10)]) # 主进程合并结果 for key, val in results: local_dict[key].append(val) for k, v in local_dict.items(): print(k, v) if __name__ == "__main__": func()
这种方式效率最高,适合不需要实时共享数据的场景。
内容的提问来源于stack exchange,提问作者Estif
相关产品推荐
相关产品推荐

