Python多线程/多进程:为每个任务添加独立进度条
并行处理多文件实现独立进度条
问题背景
并行处理多个文件时,需要为每个文件显示独立进度条,但结合并行框架与进度条工具时出现以下问题:
- 串行执行时进度条正常显示
ThreadPoolExecutor:进度条可独立推进,但最终显示混乱ProcessPoolExecutor/multiprocessing.Pool:进度条互相覆盖,无法正常显示
已测试工具:tqdm、alive_progress;并行框架:ThreadPoolExecutor、ProcessPoolExecutor、multiprocessing.Pool
解决方案
1. 线程池(ThreadPoolExecutor)场景
线程共享同一标准输出(stdout),进度条混乱源于多线程同时写入输出流的竞态问题。通过tqdm的position参数固定每个进度条位置,并使用锁同步输出操作即可解决。
修改后的代码
from bitstring import ConstBitStream import os from tqdm.auto import tqdm from concurrent.futures import ThreadPoolExecutor import threading # 全局锁,用于同步进度条输出 tqdm_lock = threading.Lock() def read_file(input): filename, limit_packet, verbose, position = input with open(filename, 'rb') as file: b = ConstBitStream(file) tot_len = len(b) # 用position固定进度条位置,lock参数传入全局锁避免输出冲突 with tqdm(total=tot_len, desc=f'Read {os.path.basename(filename)}', bar_format='{l_bar}{bar:24}{r_bar}{bar:-24b}', position=position, lock=tqdm_lock) as pbar: while b.pos < tot_len: initial_pos = b.pos # (...do stuff and advance b.pos...) pbar.update(b.pos - initial_pos) # 为每个任务分配唯一的position值,保证进度条不重叠 tasks = [(filename, limit_packet, verbose, idx) for idx, filename in enumerate(filenames)] with ThreadPoolExecutor() as p: p.map(read_file, tasks)
2. 进程池(ProcessPoolExecutor/multiprocessing.Pool)场景
进程拥有独立stdout,子进程直接创建进度条会导致输出冲突。解决方案是主进程统一管理所有进度条,子进程仅返回进度更新数据,由主进程更新对应进度条。
基于ProcessPoolExecutor的实现
from bitstring import ConstBitStream import os from tqdm.auto import tqdm from concurrent.futures import ProcessPoolExecutor def process_file(input): filename, limit_packet, verbose = input with open(filename, 'rb') as file: b = ConstBitStream(file) tot_len = len(b) while b.pos < tot_len: initial_pos = b.pos # (...do stuff and advance b.pos...) # 返回当前文件标识和进度增量 yield os.path.basename(filename), b.pos - initial_pos # 任务完成时返回结束标记 yield os.path.basename(filename), None if __name__ == '__main__': # 初始化所有文件的进度条,用字典存储便于索引 progress_bars = {} for filename in filenames: basename = os.path.basename(filename) with open(filename, 'rb') as f: tot_len = len(ConstBitStream(f)) progress_bars[basename] = tqdm(total=tot_len, desc=f'Read {basename}', bar_format='{l_bar}{bar:24}{r_bar}{bar:-24b}') # 提交任务并处理进度更新 with ProcessPoolExecutor() as executor: futures = [executor.submit(process_file, (fn, limit_packet, verbose)) for fn in filenames] for future in futures: for basename, delta in future.result(): if delta is None: # 任务完成,关闭对应进度条 progress_bars[basename].close() else: # 更新进度条 progress_bars[basename].update(delta)
关于alive_progress的适配
若偏好alive_progress,线程场景下可通过锁同步输出;进程场景下逻辑与上述一致,由主进程管理进度条、子进程返回进度数据。但alive_progress的多进程支持相对繁琐,优先推荐tqdm处理多进度条并行场景。
内容的提问来源于stack exchange,提问作者Guiux
相关产品推荐
相关产品推荐

