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

如何确认concurrent.futures.ThreadPoolExecutor是否真的并行运行?

关于Python 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()

我的问题:

  1. 如何判断我的Python代码是否真的在并行运行?
  2. 如果没有并行,该如何实现并行?

注:我希望学习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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 03:22:53