pandarallel中parallel_apply如何按批次处理行并保留并行能力
问题描述
对pandas DataFrame对象df应用自定义处理函数时,已引入pandarallel库实现流程并行化,但需要满足以下需求:每次调用func_do函数时传入N行数据,方便函数内部利用向量化运算提升处理效率。现有实现逐行调用func_do,需要调整代码,实现单次函数调用处理一个批次的行数据,同时保留原有并行处理能力。
现有参考代码:
def fun_do(value_col): return do(value_col) df['processed_col'] = df.parallel_apply(lambda row: fun_do(row['col']), axis=1)
实现方法
核心逻辑是放弃逐行axis=1的并行apply模式,改为先按指定批次大小切分DataFrame为多个子分片,以分片为单位做并行计算,每个分片整体传入处理函数走向量化逻辑,最后拼接结果即可,具体步骤如下:
- 调整自定义函数
fun_do的入参,使其直接接收一批数据的Series对象,内部直接跑向量化运算逻辑,返回和输入批次行数等长的结果序列。 - 按设定的批次大小N,按原数据顺序切分DataFrame为若干个行数不超过N的子分片。
- 对分片列表做并行映射,每个分片作为整体传入
fun_do完成批量计算,避免逐行调用的开销。 - 按原有顺序拼接所有分片的处理结果,赋值给DataFrame的目标列。
参考实现代码:
from pandarallel import pandarallel import pandas as pd # 若已完成pandarallel初始化可跳过这行 pandarallel.initialize() # 自定义单批次处理的行数,可根据内存大小、向量化逻辑效率调试调整 BATCH_N = 1000 def fun_do(batch_col: pd.Series): # 内部直接写向量化处理逻辑即可,入参是整批数据而非单行值 return do(batch_col) # 按批次大小切分原DataFrame chunks = [df.iloc[i:i+BATCH_N] for i in range(0, len(df), BATCH_N)] # 并行处理所有分片 processed_result = pd.Series(chunks).parallel_apply( lambda chunk: fun_do(chunk['col']) ) # 拼接结果赋值回原DataFrame df['processed_col'] = pd.concat(processed_result.tolist(), ignore_index=False)
注意:如果
do函数本身支持pandas/numpy向量化运算,这种实现相比逐行并行apply会有非常明显的性能提升,既保留了多进程并行能力,又省去了逐行传参、逐行调用函数的额外开销,同时能充分发挥向量化运算的效率优势。批次大小不要设置过大避免单分片内存溢出,也不要设置过小抵消向量化收益,一般在1000~10000区间调试即可。
内容的提问来源于stack exchange,提问作者H.H
相关产品推荐
相关产品推荐

