如何将asyncio异步预测结果合并至Pandas分组DataFrame(等效同步apply)
实现异步调用结果与原DataFrame的自动合并
要达到和同步groupby.apply一致的合并效果,核心是把每个异步返回的结果与对应分组的原索引绑定,再合并所有结果。具体步骤如下:
1. 分组并保存分组信息
先将分组转换为可遍历的列表,同时记录每个分组的原索引(用于后续对齐):
# 转换分组为列表,方便遍历 sales_groups = list(sales_df.groupby(level=1)) tasks = [] group_indices = [] for group_id, group_data in sales_groups: # 收集异步任务(假设get_forecast是async函数) tasks.append(get_forecast(group_data)) # 保存该分组的原MultiIndex(date, id) group_indices.append(group_data.index)
2. 异步执行并处理结果
执行所有异步任务后,将每个结果的索引替换为对应分组的原索引,确保和原DataFrame的索引对齐:
# 异步获取所有结果 results = await asyncio.gather(*tasks) # 处理每个结果,对齐索引 processed_results = [] for res, idx in zip(results, group_indices): # 针对Series/DataFrame分别处理索引 if isinstance(res, (pd.Series, pd.DataFrame)): res.index = idx processed_results.append(res) # 合并所有处理后的结果,结构与同步apply输出一致 forecasts = pd.concat(processed_results)
3. 合并至原DataFrame
此时forecasts的结构和同步groupby.apply的输出完全一致,可以直接和原DataFrame合并:
# 按索引合并,新增预测相关列 sales_df = sales_df.join(forecasts)
关键注意事项
- 确保
get_forecast返回的结果长度与对应分组的行数一致,否则索引对齐会报错。 - 如果返回的是多指标(预测值、拟合值、误差等)的DataFrame,
concat后会自动保留列名,合并后原DataFrame会直接新增这些列,和同步逻辑完全匹配。
内容的提问来源于stack exchange,提问作者mal
相关产品推荐
相关产品推荐

