如何将列表迭代索引传入multiprocessing pool调用的处理函数?
多进程池处理分块DataFrame时获取迭代索引的方法
直接在构造任务参数时为每个子DataFrame绑定对应顺序索引,一并传入进程池即可,该方案无额外性能损耗,索引顺序和单进程遍历结果完全一致。
代码修改示例
主进程参数构造部分
import numpy as np import pandas as pd import multiprocessing from multiprocessing import Pool from itertools import repeat df = pd.read_csv(input_file, encoding='utf8') dfs = np.split(df, [chunk_size]) process_pool = Pool(multiprocessing.cpu_count()) # 新增索引作为第一个传参,按顺序匹配每个子DataFrame process_pool.starmap(process_df, zip(range(len(dfs)), dfs, repeat(data_file), repeat(data_path)))
处理函数调整
# 入参新增idx接收索引 def process_df(idx, df, data_file, data_path): ... # 直接使用传入的idx生成文件名即可 output_file_name = data_path + modified_data_file + str(idx) + '.csv'
内容的提问来源于stack exchange,提问作者marlon
相关产品推荐
相关产品推荐

