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

如何利用多数据块并行执行相同处理流程?

并行处理DataFrame分块的实现方案

首先你需要先定义单个数据块的处理函数,比如针对每个块里的文件名执行读取、分析等操作,示例如下:

import pandas as pd
import os

def process_chunk(chunk_df, wb_path):
    # 这里写你对单个数据块的处理逻辑,比如读取每个文件并处理
    results = []
    for filename in chunk_df['Wk_nm']:
        file_path = os.path.join(wb_path, filename)
        # 示例:读取Excel文件(根据你的实际文件类型调整)
        df_file = pd.read_excel(file_path)
        # 这里添加你的处理逻辑,比如计算统计值、数据清洗等
        chunk_result = {
            'filename': filename,
            'row_count': len(df_file),
            'mean_value': df_file['某列'].mean() if '某列' in df_file.columns else None
        }
        results.append(chunk_result)
    # 返回该块的处理结果
    return pd.DataFrame(results)

接下来提供几种常用的并行实现方式:


方法1:使用concurrent.futures.ProcessPoolExecutor(推荐,语法简洁)

from concurrent.futures import ProcessPoolExecutor

if __name__ == '__main__':
    wb_path = "你的目录路径"
    # 你的分块代码(已给出)
    entries = os.listdir(wb_path)
    df = pd.DataFrame(entries, columns=['Wk_nm'])
    n = 50
    list_df = [df[i:i+n] for i in range(0, len(df), n)]
    
    # 并行处理
    with ProcessPoolExecutor() as executor:
        # 提交所有分块的处理任务,每个任务传入分块和目录路径
        futures = [executor.submit(process_chunk, chunk, wb_path) for chunk in list_df]
        # 收集所有结果并合并
        all_results = pd.concat([future.result() for future in futures], ignore_index=True)
    
    # 查看最终合并后的结果
    print(all_results)

方法2:使用multiprocessing.Pool(传统多进程方式)

from multiprocessing import Pool

if __name__ == '__main__':
    wb_path = "你的目录路径"
    # 你的分块代码
    entries = os.listdir(wb_path)
    df = pd.DataFrame(entries, columns=['Wk_nm'])
    n = 50
    list_df = [df[i:i+n] for i in range(0, len(df), n)]
    
    # 初始化进程池(默认使用CPU核心数)
    with Pool() as pool:
        # 使用starmap传递多个参数(分块和目录路径)
        results = pool.starmap(process_chunk, [(chunk, wb_path) for chunk in list_df])
        # 合并结果
        all_results = pd.concat(results, ignore_index=True)
    
    print(all_results)

方法3:使用joblib(适合批量任务,支持进度条)

如果你需要查看处理进度,可以用joblib,需要先安装:pip install joblib

from joblib import Parallel, delayed

if __name__ == '__main__':
    wb_path = "你的目录路径"
    # 你的分块代码
    entries = os.listdir(wb_path)
    df = pd.DataFrame(entries, columns=['Wk_nm'])
    n = 50
    list_df = [df[i:i+n] for i in range(0, len(df), n)]
    
    # 并行处理,n_jobs指定进程数,verbose显示进度
    results = Parallel(n_jobs=-1, verbose=10)(
        delayed(process_chunk)(chunk, wb_path) for chunk in list_df
    )
    all_results = pd.concat(results, ignore_index=True)
    
    print(all_results)

注意事项

  • 必须把主逻辑放在if __name__ == '__main__':下,避免多进程启动时重复执行代码
  • 处理函数process_chunk要能被序列化(pickle),所以不要在函数内部定义嵌套函数或使用无法序列化的对象
  • 根据你的CPU核心数调整进程数,默认会使用全部核心,也可以手动指定(比如ProcessPoolExecutor(max_workers=4))

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 18:12:52