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

Python ThreadPoolExecutor未按预期实现并发运行问题排查

问题描述

在Python中使用ThreadPoolExecutor时,代码运行无错误且输出正确,但未实现预期的并发执行。单个最长任务耗时约4小时,批处理总耗时却达10小时,且机器有足够核心可同时运行所有任务。

代码逻辑如下:

import a 
import b 
import c 
import d 
import e 
from f import g
from t import q

funcs = [a.func1, b.func2, [g.func3, g.func4, e.func5], c.func8, d.func6, [q.func9, q.func10]]


def check_and_run(funcs, out_text: str = "") -> str:
    if type(funcs) == list:
        for func in funcs:
            out_text = check_and_run(func, out_text=out_text)
        return out_text
    else:
        func = funcs
        try:
            out_text = func(out_text=out_text)
            return out_text
        except AttributeError as attr:
            if str(attr) == "'functools.partial' object has no attribute '__name__'":
                return f"{out_text}\nERROR - check command line - {str(func)}"
        except Exception as e:
            return f"{out_text}\nERROR:\t{func.__name__}, \t{str(e)}"

with ThreadPoolExecutor(max_workers=default_num_cores) as executor:
   for routine_or_list_of_routines in list_of_routines:
        out_result = executor.submit(check_and_run, funcs)
        temp = out_result.result()
        temp = "".join(i for i in temp) if type(temp) == list else temp
        out_text = f"{out_text}{temp}"

场景说明:

  • 有一个列表list_of_routines,元素为单个待执行函数或依赖型函数列表(部分任务需依赖其他任务执行)。
  • 每个函数接收同一字符串参数,并追加运行状态信息(用于批处理结束后查看运行状态)。
问题原因分析
  • 同步阻塞调用result():循环中每次提交任务后立刻调用out_result.result(),会强制当前线程等待该任务完成才进入下一次循环,直接把并发变成了串行执行,完全没用到线程池的并发能力。
  • 线程池适用场景偏差:如果任务是CPU密集型,Python的GIL(全局解释器锁)会限制多线程的实际并行能力,此时ThreadPoolExecutor无法利用多核;如果是IO密集型任务,你的代码也因为阻塞调用没发挥线程池的作用。
  • 任务提交逻辑错误:循环里每次提交的都是check_and_run(funcs),而非routine_or_list_of_routines,等于重复执行同一个funcs列表的任务,不符合预期的任务分发逻辑。
解决方案

方案1:修复线程池并发调用逻辑

将result()的调用移到所有任务提交之后,避免阻塞循环:

from concurrent.futures import ThreadPoolExecutor

with ThreadPoolExecutor(max_workers=default_num_cores) as executor:
    # 先批量提交所有任务,收集future对象
    futures = []
    for routine in list_of_routines:
        # 注意这里传参为当前循环的routine,而非固定的funcs
        future = executor.submit(check_and_run, routine)
        futures.append(future)
    
    # 统一等待所有任务完成,处理结果
    out_text = ""
    for future in futures:
        temp = future.result()
        temp = "".join(i for i in temp) if isinstance(temp, list) else temp
        out_text += temp

方案2:CPU密集型任务改用进程池

如果任务是CPU密集型,替换为ProcessPoolExecutor以利用多核(确保函数和参数可被序列化):

from concurrent.futures import ProcessPoolExecutor

with ProcessPoolExecutor(max_workers=default_num_cores) as executor:
    futures = []
    for routine in list_of_routines:
        future = executor.submit(check_and_run, routine)
        futures.append(future)
    
    out_text = ""
    for future in futures:
        temp = future.result()
        temp = "".join(i for i in temp) if isinstance(temp, list) else temp
        out_text += temp

方案3:优化check_and_run的类型判断与异常处理

原代码类型判断不够严谨,同时优化异常中对函数名的获取逻辑:

def check_and_run(funcs, out_text: str = "") -> str:
    # 用isinstance兼容列表子类,比直接判断type更严谨
    if isinstance(funcs, list):
        for func in funcs:
            out_text = check_and_run(func, out_text=out_text)
        return out_text
    else:
        func = funcs
        try:
            out_text = func(out_text=out_text)
            return out_text
        except AttributeError as attr:
            if "'functools.partial' object has no attribute '__name__'" in str(attr):
                return f"{out_text}\nERROR - check command line - {str(func)}"
        except Exception as e:
            # 兼容无法获取__name__的情况
            func_name = getattr(func, "__name__", "unknown function")
            return f"{out_text}\nERROR:\t{func_name}, \t{str(e)}"

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 14:42:01