如何确认concurrent.futures.ThreadPoolExecutor是否真的并行运行?
我尝试用Python的concurrent.futures.ThreadPoolExecutor类并行化一个向列表追加元素的逻辑,代码能正常运行,但不确定是否真的在执行多线程操作。根据Python官方文档,Executor实例的.map(func, iterable)方法应该能以任意顺序对可迭代对象的每个元素调用函数:
map(func, *iterables, timeout=None, chunksize=1)
Similar to map(func, *iterables) except:
the iterables are collected immediately rather than lazily;
func is executed asynchronously and several calls to func may
be made concurrently.
但我的代码输出顺序和可迭代对象完全一致,这让我怀疑它并没有真正并行运行。我的Python代码如下:
from concurrent.futures import ThreadPoolExecutor import itertools def foo(xy): x, y = xy return f"{x} + {y} = {x + y}" if __name__ == "__main__": n = 8 with ThreadPoolExecutor(4) as e: prod = itertools.product(range(0, n), range(n, 2 * n)) lst = list(e.map(foo, prod)) assert len(lst) == n * n # 确保所有结果都被收集 print(lst)
它的输出顺序和itertools.product生成的迭代器完全一致。
对比Julia的@threads代码,输出是无序的,说明是并行执行的:
using Base.Threads foo(x, y)::String = "$x + $y = $(x + y)" function main() vec = String[] n = 8 tuples = Iterators.product(1:n, n+1:2n) |> collect @threads for (x, y) in tuples push!(vec, foo(x, y)) end @assert length(vec) == n * n display(vec) end main()
我的问题:
- 如何判断我的Python代码是否真的在并行运行?
- 如果没有并行,该如何实现并行?
注:我希望学习concurrent.futures.ThreadPoolExecutor类,而非使用multiprocessing.dummy.Pool。
问题解答
1. 如何判断Python代码是否在并行运行?
你的代码确实是并行运行的,只是e.map()方法会自动按输入迭代器的顺序整理返回结果,所以输出看起来是有序的——这是该方法的设计特性,和是否并行无关。
要验证并行性,可以通过以下方式:
- 添加随机延迟:在
foo函数里加入time.sleep(random.uniform(0.1, 0.5)),如果是并行执行,任务的实际完成顺序会打乱,但map返回的列表仍会保持输入顺序;如果是串行,输出顺序和任务完成顺序一致,且总耗时接近所有延迟的总和。 - 打印线程ID:导入
threading模块,在foo函数里添加print(f"当前线程ID: {threading.get_ident()}"),运行后会看到多个不同的线程ID,说明多线程在工作。 - 对比执行耗时:写一个串行版本的代码(直接用
list(map(foo, prod))),统计并行和串行的总耗时,并行版本总耗时会远小于串行(尤其适合IO密集型任务)。
2. 实现并行并得到无序结果的方法
如果想让输出结果反映实际的任务完成顺序(即无序),不要用map方法,改用submit配合as_completed:
from concurrent.futures import ThreadPoolExecutor, as_completed import itertools import time import random def foo(xy): x, y = xy # 添加随机延迟模拟真实任务场景 time.sleep(random.uniform(0.1, 0.5)) return f"{x} + {y} = {x + y}" if __name__ == "__main__": n = 8 lst = [] with ThreadPoolExecutor(4) as e: prod = itertools.product(range(0, n), range(n, 2 * n)) # 提交所有任务,获取future对象列表 futures = [e.submit(foo, item) for item in prod] # 按任务完成的先后顺序获取结果 for future in as_completed(futures): lst.append(future.result()) assert len(lst) == n * n print(lst)
这段代码的输出会是无序的,因为as_completed会在任务完成时立即返回结果,完全反映并行执行的真实顺序。
另外需要注意:Python的线程受GIL(全局解释器锁)限制,CPU密集型任务无法通过线程池实现真正的并行加速,此时需要改用ProcessPoolExecutor;但IO密集型任务(比如网络请求、文件读写)能通过线程池有效提升效率,因为线程在等待IO时会释放GIL,让其他线程执行。
内容的提问来源于stack exchange,提问作者Max

