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

如何从ThreadPoolExecutor任务中提取DataFrame并完成合并?

如何从ThreadPoolExecutor任务中提取DataFrame并完成合并?

我来帮你拆解一下问题,你遇到的两个错误其实是两个不同的问题,咱们一步步解决:

第一个错误的原因:你直接在Future对象上调用了pd.concat

wait()方法返回的task_results是Future对象的集合,不是实际的DataFrame结果。pd.concat()只认识DataFrame或Series对象,自然会报错说无法合并Future类型的对象。这一步的解决思路很明确:必须从每个Future对象里取出真正的DataFrame结果。

第二个错误的原因:你的任务函数写错了

当你尝试用[t.result() for t in task_results]时,报错“DataFrame不可调用”,这说明你提交给executor.submit()的任务逻辑有问题——你大概率是把一个已经存在的DataFrame对象当成函数来调用了。举个例子:如果你的代码里写了lambda: my_df(),但my_df是一个已经创建好的DataFrame(不是返回DataFrame的函数),那Python就会尝试调用这个DataFrame,自然触发“不可调用”的错误。

正确的解决步骤

1. 先确保任务函数能正确返回DataFrame

不管你用自定义函数还是lambda,都要保证它是返回DataFrame,而不是调用DataFrame。比如:

  • 用自定义函数(更清晰,推荐):
def my_data_processing():
    # 这里写你的实际计算逻辑,比如处理原DataFrame的某一部分
    return pd.DataFrame({'col': [1,2,3]})
  • 用lambda的正确写法:如果是返回已有的DataFrame,不要加括号;如果是动态创建,就正常写:
# 动态创建DataFrame的lambda
lambda: pd.DataFrame({'col': [1,2,3]})
# 返回已有DataFrame的lambda(注意不要加括号)
lambda: existing_df.copy()

2. 正确提取结果并合并

等所有任务完成后,遍历每个Future对象调用result()方法获取DataFrame,再用pd.concat()合并:

完整示例代码:

from concurrent.futures import ThreadPoolExecutor, wait, ALL_COMPLETED
import pandas as pd

# 替换成你实际的数据处理函数
def process_chunk():
    # 这里模拟你的DataFrame计算逻辑
    return pd.DataFrame({'id': [1,2], 'value': [10,20]})

executor = ThreadPoolExecutor(max_workers=20)

# 提交40个任务
tasks = [executor.submit(process_chunk) for _ in range(40)]

# 等待所有任务完成
done_tasks, _ = wait(tasks, return_when=ALL_COMPLETED)

# 提取所有DataFrame结果
df_results = [task.result() for task in done_tasks]

# 合并所有DataFrame
final_df = pd.concat(df_results, ignore_index=True)

这样就能顺利得到合并后的DataFrame啦!

备注:内容来源于stack exchange,提问作者wedrano de carvalho

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.22 10:29:32