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

使用multiprocessing与tqdm时,如何在进度条下方打印任务状态

问题场景

使用multiprocessing结合tqdm展示多进程任务进度时,希望在进度条下方添加一行实时显示当前正在处理的任务,但所有额外输出都会跑到进度条上方,导致界面混乱。复现代码如下:

from multiprocessing import Pool, Manager, Value
import time
import os
import tqdm
import sys

class ParallelProcessor:
    def __init__(self, shared_data):
        self.shared_data = shared_data

    def process_task(self, args):
        """Worker function: Simulates task processing and updates progress"""
        lock, progress, active_tasks, index, integer_arg = args
        pid = os.getpid()
        core_id = index % len(os.sched_getaffinity(0))
        os.sched_setaffinity(pid, {core_id})

        with lock:
            active_tasks.append(f"Task {index+1}")

        time.sleep(2)  # Simulate processing time

        with lock:
            active_tasks.remove(f"Task {index+1}")
            progress.value += 1

        return self.shared_data

    def progress_updater(self, total_tasks, progress, active_tasks):
        """Update tqdm progress bar and active task list on separate lines"""
        sys.stdout.write("\n")  # Move to the next line for active task display
        sys.stdout.flush()

        with tqdm.tqdm(total=total_tasks, desc="Processing Tasks", position=0, leave=True) as pbar:
            while pbar.n < total_tasks:
                time.sleep(0.1)  # Update interval
                pbar.n = progress.value
                pbar.refresh()

                # Move cursor down to the next line and overwrite active task display
                sys.stdout.write("\033[s")  # Save cursor position
                sys.stdout.write(f"\033[2K\rActive: {', '.join(active_tasks[:5])}")  # Clear line and print active tasks
                sys.stdout.write("\033[u")  # Restore cursor position
                sys.stdout.flush()

    def run_parallel(self, tasks, num_cores=None):
        """Runs tasks in parallel with a progress bar"""
        num_cores = num_cores or len(os.sched_getaffinity(0))
        manager = Manager()
        lock = manager.Lock()
        progress = manager.Value("i", 0)  # Shared integer for progress tracking
        active_tasks = manager.list()  # Shared list for active tasks

        # Start progress updater in the main process
        from threading import Thread
        progress_thread = Thread(target=self.progress_updater, args=(len(tasks), progress, active_tasks))
        progress_thread.start()

        # Prepare task arguments
        task_args = [(lock, progress, active_tasks, idx, val) for idx, val in enumerate(tasks)]

        # Run parallel tasks
        with Pool(num_cores) as pool:
            results = pool.map(self.process_task, task_args)

        # Ensure progress bar finishes
        progress_thread.join()
        print("\n")  # Move to the next line after processing

        return results


if __name__ == "__main__":
    processor = ParallelProcessor(shared_data=10)
    processor.run_parallel(tasks=range(40), num_cores=4)
解决方案

核心思路是利用终端光标控制指令固定进度条和任务显示行的位置,避免手动光标操作的混乱,同时确保线程安全的任务状态更新:

  1. 将进度条固定在终端第0行,任务显示行固定在第1行
  2. 用光标移动指令在两行之间切换,更新任务列表前先清除当前行内容
  3. 增加异常处理避免多进程竞态导致的任务移除错误

修改后的完整代码:

from multiprocessing import Pool, Manager, Value
import time
import os
import tqdm
import sys

class ParallelProcessor:
    def __init__(self, shared_data):
        self.shared_data = shared_data

    def process_task(self, args):
        lock, progress, active_tasks, index, integer_arg = args
        pid = os.getpid()
        core_id = index % len(os.sched_getaffinity(0))
        os.sched_setaffinity(pid, {core_id})

        with lock:
            active_tasks.append(f"Task {index+1}")

        time.sleep(0.5)  # 缩短模拟时间方便测试

        with lock:
            try:
                active_tasks.remove(f"Task {index+1}")
            except ValueError:
                pass  # 避免任务重复移除的异常
            progress.value += 1

        return self.shared_data

    def progress_updater(self, total_tasks, progress, active_tasks):
        # 初始化进度条在第0行
        with tqdm.tqdm(total=total_tasks, desc="Processing Tasks", position=0, leave=True) as pbar:
            # 先清除第1行的初始空白内容
            sys.stdout.write("\033[1B\033[2K\r\033[1A")
            sys.stdout.flush()
            
            while pbar.n < total_tasks:
                time.sleep(0.1)
                pbar.n = progress.value
                pbar.refresh()

                # 光标下移到第1行,清除该行并写入活跃任务
                sys.stdout.write("\033[1B")
                sys.stdout.write(f"\033[2K\rActive: {', '.join(active_tasks) if active_tasks else 'None'}")
                # 光标移回进度条行
                sys.stdout.write("\033[1A")
                sys.stdout.flush()

        # 处理完成后清除活跃任务行
        sys.stdout.write("\033[1B\033[2K\r")
        sys.stdout.flush()

    def run_parallel(self, tasks, num_cores=None):
        num_cores = num_cores or len(os.sched_getaffinity(0))
        manager = Manager()
        lock = manager.Lock()
        progress = manager.Value("i", 0)
        active_tasks = manager.list()

        from threading import Thread
        progress_thread = Thread(target=self.progress_updater, args=(len(tasks), progress, active_tasks))
        progress_thread.start()

        task_args = [(lock, progress, active_tasks, idx, val) for idx, val in enumerate(tasks)]

        with Pool(num_cores) as pool:
            results = pool.map(self.process_task, task_args)

        progress_thread.join()
        print("\nProcessing complete!")

        return results


if __name__ == "__main__":
    processor = ParallelProcessor(shared_data=10)
    processor.run_parallel(tasks=range(40), num_cores=4)
关键改动说明
  • 精准光标控制:用\033[1B(下移1行)和\033[1A(上移1行)固定光标在进度条和任务行之间切换
  • 行内容清理:每次更新任务列表前用\033[2K清除当前行,避免残留旧内容
  • 竞态处理:在任务移除逻辑中增加try-except,防止多进程同时操作导致的ValueError
  • 界面收尾:任务完成后清除任务显示行,保证终端界面整洁

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 09:40:56