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'
性能瓶颈分析
- DataFrame无索引的线性查询:每次迭代执行
df[df[id_name] == item[id_name]],对百万级DataFrame是O(n)线性扫描,多次调用后总复杂度变为O(m*n)(m为实体数量),完全浪费了Pandas的索引优势。 - 跨进程共享对象的锁开销:
multiprocessing.Manager的共享机制会触发严重的进程同步锁阻塞,直接抵消多进程的并行收益。 - 低效的字符串拼接与共享写入:用
+=拼接字符串本身效率极低,跨进程写入共享的obj['json'][name]还会引发锁竞争,进一步拖慢速度。 - 频繁创建销毁进程池:对每个实体类型新建进程池,进程初始化与销毁的累加开销极大。
- 重复的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
相关产品推荐
相关产品推荐

