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

Python百万级字典迭代追加的性能优化求助

批量处理实体数据的性能优化问题

问题背景

处理数百万条字典数据(来自多个单文件100万行的数据源),最终生成Elasticsearch批量导入用的JSON文件。核心需求是将“实体”数据与“子实体”中的对应地址关联,但当前关联函数的单次迭代耗时极长,优化后速度仍未达标。

数据结构说明

  • 实体(如Persons):包含id、uniqueId、name、postalList、emailList字段,示例:
    {id: 0, uniqueId: 'Z5ER1ZE5', name: 'John DOE', postalList: [], emailList: []}
    
  • 子实体(不同类型地址):包含关联实体的唯一ID及对应地址,示例:
    {personUniqueId: 'Z5ER1ZE5', 'Email': 'john.doe@gmail.com'}
    

现有实现代码

多进程初始化代码

## 使用Manager实现跨进程对象共享与更新
manager = multiprocessing.Manager()
main_obj = manager.dict({
   'dataframes': manager.dict(),
   'dicts': manager.dict(),
   'json': manager.dict()
})

## 多进程调用示例
pool = multiprocessing.Pool()
result = pool.map(partial(func, obj=main_obj), data_list_to_iterate)
pool.close()
pool.join()

参考关联字典

sub_associations = {
   'Persons_emlAddr': 'postalList',
   'Persons_pstlAddr': 'emailList'
}
identifiers = {
   'Animals': 'uniqueAnimalId',
   'Persons': 'uniquePersonId',
   'Persons_emlAddr': 'uniquePersonId',
   'Persons_pstlAddr': 'uniquePersonId'
}

子实体关联与JSON生成逻辑

for key in list(main_obj['dicts'].keys()):
        main_obj['json'][key] = ''
        with mp.Pool() as stringify_pool:
            res_stringify = stringify_pool.map(partial(convert_to_json, obj=main_obj, name=key), main_obj['dicts'][key]['records'])
            stringify_pool.close()
            stringify_pool.join()

核心关联函数convert_to_json

def convert_to_json(item, obj, name):
    global sub_associations
    global identifiers
    
    dump = ''
    
    subs = [val for val in sub_associations.keys() if val.startswith(name)]
    if subs:
        for sub in subs:
            df = obj['dataframes'][sub]
            
            id_name = identifiers[name]
            sub_items = df[df[id_name] == item[id_name]].to_dict('records')
            
            if sub_items:
                item[sub_associations[sub]] = sub_items
            else:
                item[sub_associations[sub]] = []
    
    index = {
        "index": {
            "_index": name,
            "_id": item[identifiers[name]]
        }
    }
            
    dump += f'{json.dumps(index)}\n'
    dump += f'{json.dumps(item)}\n'
    
    obj['json'][name] += dump
    
    return 'Done'

性能瓶颈分析

  1. DataFrame无索引的线性查询:每次迭代执行df[df[id_name] == item[id_name]],对百万级DataFrame是O(n)线性扫描,多次调用后总复杂度变为O(m*n)(m为实体数量),完全浪费了Pandas的索引优势。
  2. 跨进程共享对象的锁开销:multiprocessing.Manager的共享机制会触发严重的进程同步锁阻塞,直接抵消多进程的并行收益。
  3. 低效的字符串拼接与共享写入:用+=拼接字符串本身效率极低,跨进程写入共享的obj['json'][name]还会引发锁竞争,进一步拖慢速度。
  4. 频繁创建销毁进程池:对每个实体类型新建进程池,进程初始化与销毁的累加开销极大。
  5. 重复的DataFrame转字典操作:每次查询后都执行to_dict('records'),重复的序列化操作带来额外性能损耗。

优化方案

1. 预对子实体数据建立映射索引

提前按关联ID分组子实体数据,生成字典映射,彻底避免每次查询DataFrame:

# 预处理子实体数据,生成id到子实体列表的映射
sub_maps = {}
for sub_name, target_field in sub_associations.items():
    id_col = identifiers[sub_name]
    df = obj['dataframes'][sub_name]
    # 按id分组并转成字典列表
    sub_maps[sub_name] = df.groupby(id_col).apply(lambda x: x.to_dict('records')).to_dict()

2. 移除共享对象,改为进程传递数据

不再用Manager共享数据,把预处理好的sub_maps、identifiers、sub_associations直接传给子进程,每个进程独立处理任务并返回结果片段:

# 全局创建一次进程池,复用所有任务
with multiprocessing.Pool() as pool:
    for key in list(main_obj['dicts'].keys()):
        records = main_obj['dicts'][key]['records']
        # 传递预处理好的映射而非共享对象
        args = [(item, key, sub_maps, identifiers, sub_associations) for item in records]
        # 获取所有结果片段并汇总
        json_fragments = pool.map(process_item, args)
        main_obj['json'][key] = ''.join(json_fragments)

3. 优化核心处理函数为无共享模式

修改后的函数不再依赖共享对象,直接使用传入的映射数据,并用列表拼接提升字符串效率:

def process_item(args):
    item, name, sub_maps, identifiers, sub_associations = args
    dump_parts = []
    
    id_name = identifiers[name]
    item_id = item[id_name]
    
    subs = [val for val in sub_associations.keys() if val.startswith(name)]
    for sub in subs:
        target_field = sub_associations[sub]
        # 直接从预生成的映射中取值
        item[target_field] = sub_maps[sub].get(item_id, [])
    
    # 生成index行
    index = {"index": {"_index": name, "_id": item_id}}
    dump_parts.append(json.dumps(index))
    dump_parts.append(json.dumps(item))
    
    return '\n'.join(dump_parts) + '\n'

4. 其他优化点

  • 避免全局变量:把关联字典作为参数传递,减少进程间的变量复制开销。
  • 批量处理优化:如果内存允许,考虑用Pandas批量处理实体与子实体的关联,再统一生成JSON。
  • 内存压力缓解:若数据量超出内存,用Dask分块处理数据,避免一次性加载所有DataFrame。

内容的提问来源于stack exchange,提问作者Oraluka

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 10:55:20