使用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)
解决方案
核心思路是利用终端光标控制指令固定进度条和任务显示行的位置,避免手动光标操作的混乱,同时确保线程安全的任务状态更新:
- 将进度条固定在终端第0行,任务显示行固定在第1行
- 用光标移动指令在两行之间切换,更新任务列表前先清除当前行内容
- 增加异常处理避免多进程竞态导致的任务移除错误
修改后的完整代码:
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
相关产品推荐
相关产品推荐

