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
相关产品推荐
相关产品推荐

