如何用Python多进程处理简单列表以并行执行数据处理任务?
多进程处理唯一Key列表的实现方案
基础实现(基于你现有代码)
直接使用multiprocessing.Pool的imap_unordered或map方法,就能把vidList中的每个元素传入多进程任务,同时收集返回结果。
完整代码示例
import multiprocessing import functools # 你已定义的taska函数 def taska(id, data, key, date): cdata = data[data[key]==id] return vt.make_breakpts(data=cdata, key=key, date=date) if __name__ == '__main__': # 固定不变参数(你的原有代码) partialA = functools.partial(taska, data=pgdf, key=VID, date=PDATE) resultList = [] vidList = list(pgdf['VESSEL_ID'].unique()) # 初始化进程池,默认使用CPU核心数 with multiprocessing.Pool() as pool: # 方案1:用imap_unordered,返回结果顺序不固定(速度略快) for res in pool.imap_unordered(partialA, vidList): resultList.append(res) # 方案2:用map,返回结果顺序与vidList完全对应 # resultList = pool.map(partialA, vidList)
内存优化方案(针对大DataFrame)
如果你的pgdf数据量很大,每个进程复制完整DataFrame会占用大量内存。可以先按Key分组,直接传递子DataFrame给进程:
优化后代码
import multiprocessing import functools # 修改taska函数,直接接收子DataFrame def taska(id, cdata, key, date): return vt.make_breakpts(data=cdata, key=key, date=date) if __name__ == '__main__': # 先按Key分组,避免重复筛选 grouped_data = [(vid, group) for vid, group in pgdf.groupby(VID)] partialA = functools.partial(taska, key=VID, date=PDATE) resultList = [] with multiprocessing.Pool() as pool: # 用starmap传递元组中的多个参数 resultList = pool.starmap(partialA, grouped_data)
关键注意事项
- 必须把多进程启动代码放在
if __name__ == '__main__':块内,这是Windows系统多进程启动的强制要求,也能避免Linux/macOS下的重复初始化问题。 imap_unordered适合不关心结果顺序的场景,返回速度更快;map会严格保持输入列表的顺序。- 若
vt.make_breakpts的返回结果是DataFrame,最终resultList会是DataFrame列表,可通过pd.concat(resultList)合并成完整结果。
内容的提问来源于stack exchange,提问作者MrKingsley
相关产品推荐
相关产品推荐

