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

Python多进程Queue如何保证LOESS回归结果的持久顺序?

解决多进程LOESS回归后DataFrame顺序混乱的问题

你的核心问题在于用队列收集结果时,进程完成顺序不确定导致数据顺序错乱。下面提供两种可靠的解决方法,优先推荐第一种:

方法一:利用starmap的天然顺序特性(推荐)

Pool.starmap会严格按照你传入的任务列表顺序返回结果,不管进程实际完成时间早晚,完全不需要队列就能保证顺序。

修改步骤:

  1. 调整process_data函数:去掉队列参数,直接返回处理后的DataFrame
def process_data(dataframe, i, chunk_size, fraction):
    start_frame = chunk_size * i
    end_frame = min(chunk_size * (i + 1), len(dataframe))
    print(f'{start_frame}, {end_frame}') # 调试用
    return calculate_loess_on_subset(dataframe[start_frame:end_frame], chunk_size, fraction, i)
  1. 重构smooth_data_mp函数:用starmap的返回值构建结果,自动保证顺序
def smooth_data_mp(data_frame):
    num_processes = 8
    chunk_size = 125000
    fraction = 125 / chunk_size
    print(data_frame.head())
    
    # 计算总chunk数,确保覆盖所有数据(包括最后不足chunk_size的部分)
    num_chunks = (len(data_frame) + chunk_size - 1) // chunk_size
    # 生成任务列表,按原数据顺序分配chunk索引
    tasks = [(data_frame, i, chunk_size, fraction) for i in range(num_chunks)]
    
    with Pool(processes=num_processes) as pool:
        # starmap按任务顺序返回结果列表
        result_list = pool.starmap(process_data, tasks)
    
    # 按顺序拼接所有DataFrame
    return pd.concat(result_list, ignore_index=True)

方法二:给队列结果添加索引标记(兼容原队列写法)

如果一定要保留队列,可以让每个进程把chunk索引i和结果DataFrame一起存入队列,收集后按索引排序再拼接:

修改步骤:

  1. 更新process_data函数:存入队列时附带索引i
def process_data(dataframe, i, chunk_size, fraction, result_queue):
    start_frame = chunk_size * i
    end_frame = min(chunk_size * (i + 1), len(dataframe))
    print(f'{start_frame}, {end_frame}') # 调试用
    new_data_frame = calculate_loess_on_subset(dataframe[start_frame:end_frame], chunk_size, fraction, i)
    result_queue.put( (i, new_data_frame) ) # 把索引和结果绑定
  1. 调整结果收集逻辑:按索引排序后再拼接
def smooth_data_mp(data_frame):
    num_processes = 8
    chunk_size = 125000
    fraction = 125 / chunk_size
    print(data_frame.head())
    
    result_queue = Manager().Queue()
    num_chunks = (len(data_frame) + chunk_size - 1) // chunk_size
    tasks = [(data_frame, i, chunk_size, fraction, result_queue) for i in range(num_chunks)]
    
    with Pool(processes=num_processes) as pool:
        pool.starmap(process_data, tasks)
    
    # 收集所有带索引的结果
    result_list = []
    while not result_queue.empty():
        result_list.append(result_queue.get())
    
    # 按chunk索引排序,恢复原数据顺序
    result_list.sort(key=lambda x: x[0])
    # 提取DataFrame并拼接
    final_df = pd.concat([df for _, df in result_list], ignore_index=True)
    return final_df

关键说明:

  • 原代码中range(len(data_frame)//chunk_size)会漏掉最后不足一个chunk_size的数据块,修改为(len(data_frame) + chunk_size -1)//chunk_size可以确保所有数据都被处理。
  • 方法一利用starmap的内置顺序保证,代码更简洁高效,不需要额外的排序或队列操作,是最优解。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 23:35:26