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

如何在ProcessPoolExecutor中用tqdm避免进度条错位?

解决ProcessPoolExecutor中tqdm进度条错位的Python方案

多进程环境下,每个子进程独立输出tqdm进度条会导致终端输出混乱——多个进程同时写入stdout,进度条的控制字符互相干扰,最终出现错位、重叠问题。下面是符合Python风格的解决方案:通过进程间队列传递进度信息,由父进程统一管理并更新所有进度条,彻底避免多进程输出冲突。

实现思路

  1. 用multiprocessing.Queue作为子进程向父进程传递进度消息的通道
  2. 父进程为每个任务创建独立的outer/inner进度条,统一维护
  3. 启动单独线程监听队列,接收子进程的进度消息并更新对应进度条
  4. 子进程执行任务时,将进度数据(任务ID、进度类型、当前进度)发送到队列,不再直接输出tqdm

完整代码示例

from tqdm import tqdm
import time
import random
from concurrent.futures import ProcessPoolExecutor
import multiprocessing
import threading

def main():
    # 创建进程间通信队列
    progress_queue = multiprocessing.Queue()
    task_count = 10
    # 存储每个任务的进度条:key为任务ID,值包含outer和inner进度条
    progress_bars = {}

    # 初始化所有任务的进度条,分配固定终端行位置
    for task_id in range(task_count):
        outer_bar = tqdm(total=500, desc=f"{task_id}_outer", position=task_id*2)
        inner_bar = tqdm(total=500, desc=f"{task_id}_inner", position=task_id*2 + 1)
        progress_bars[task_id] = {"outer": outer_bar, "inner": inner_bar}

    # 监听队列的线程函数:处理进度更新
    def listen_progress():
        while True:
            msg = progress_queue.get()
            if msg is None:  # 结束信号,终止线程
                break
            task_id, bar_type, current = msg
            bar = progress_bars[task_id][bar_type]
            # 计算需要更新的步数,确保进度条准确跳转
            bar.update(current - bar.n)
            # 进度完成时关闭进度条
            if bar.n >= bar.total:
                bar.close()

    # 启动监听线程
    listener_thread = threading.Thread(target=listen_progress)
    listener_thread.start()

    # 提交任务到进程池
    with ProcessPoolExecutor() as pool:
        futures = [pool.submit(my_function, task_id, progress_queue) for task_id in range(task_count)]
        # 等待所有任务完成
        for future in futures:
            future.result()

    # 发送结束信号,等待监听线程退出
    progress_queue.put(None)
    listener_thread.join()

    # 清理剩余进度条
    for bars in progress_bars.values():
        bars["outer"].close()
        bars["inner"].close()

def my_function(fn_number, progress_queue):
    for i in range(500):
        # 发送outer循环的进度更新
        progress_queue.put((fn_number, "outer", i+1))
        inner_progress = 0
        for j in range(500):
            inner_progress += 1
            if random.random() > 0.9:
                break
            time.sleep(0.001)  # 缩短sleep时间加快测试
        # 发送inner循环的最终进度更新
        progress_queue.put((fn_number, "inner", inner_progress))
        # 重置inner进度条(下一次outer循环重新开始)
        progress_queue.put((fn_number, "inner", 0))
    return fn_number

if __name__=='__main__':
    main()

关键细节说明

  • 队列通信:子进程通过progress_queue发送进度元组(任务ID, 进度条类型, 当前进度),父进程监听线程接收后更新对应进度条,保证进程间通信安全
  • 进度条定位:用position参数为每个任务的outer/inner进度条分配固定终端行,避免不同任务的进度条互相覆盖
  • 结束信号:所有任务完成后,父进程发送None到队列,通知监听线程优雅退出
  • 进度更新逻辑:通过bar.update(current - bar.n)计算需要更新的步数,处理inner循环中途中断的场景,确保进度条状态准确

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 23:43:16