如何使用 concurrent.futures map 与 tqdm 进度条

问题:

你有一个 concurrent.futures 执行器,例如

concurrent_map_tqdm_example.py
import concurrent.futures

executor = concurrent.futures.ThreadPoolExecutor(64)

使用此执行器,你想并行地将函数 map 到可迭代对象上(例如并行下载 HTTP 页面)。

为了辅助交互式执行,你想使用 tqdm 提供进度条,显示 future 的分数

解决方案

你可以使用此函数:

concurrent_map_tqdm_example_2.py
from tqdm import tqdm
import concurrent.futures

def tqdm_parallel_map(executor, fn, *iterables, **kwargs):
    """
    等同于 executor.map(fn, *iterables),
    但显示基于 tqdm 的进度条。

    不支持 timeout 或 chunksize,因为内部使用 executor.submit

    **kwargs 传递给 tqdm。
    """
    futures_list = []
    for iterable in iterables:
        futures_list += [executor.submit(fn, i) for i in iterable]
    for f in tqdm(concurrent.futures.as_completed(futures_list), total=len(futures_list), **kwargs):
        yield f.result()

注意内部使用 executor.submit(),而不是 executor.map(),因为无法对 executor.map() 返回的迭代器调用 concurrent.futures.as_completed()

如果你对实际值不感兴趣的用法示例:

concurrent_map_tqdm_example_3.py
import concurrent.futures
executor = concurrent.futures.ThreadPoolExecutor(64)

# 并行运行 my_func,参数范围从 1 到 10000
for _ in tqdm_parallel_map(executor, lambda i: my_func(i), range(1, 10000)):
    pass

如果你关心 my_func() 的返回值,请改用此代码片段:

concurrent_map_tqdm_example_4.py
import concurrent.futures
executor = concurrent.futures.ThreadPoolExecutor(64)

# 并行运行 my_func,参数范围从 1 到 10000
for result in tqdm_parallel_map(executor, lambda i: my_func(i), range(1, 10000)):
    # result 是 my_func() 的返回值。
    # 注意:result 不一定与输入 Iterable 的顺序相同!
    # 哪个并行执行的 my_func() 先完成将先打印!
    print(result)

注意:executor.map() 相反,此函数不会按与输入相同的顺序生成参数。


Check out similar posts by category: Python