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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 15:56:02