如何结合ProcessPoolExecutor与tqdm实现平滑进度展示
多进程下实现tqdm进度条平滑更新
切换到ProcessPoolExecutor后,由于进程间内存不共享,直接传递tqdm实例会导致程序崩溃,目前只能在子任务完成后批量更新进度条,导致进度显示跳跃。以下是两种可行的平滑更新方案:
方案一:使用共享内存计数器
通过multiprocessing.Value创建跨进程共享的计数器,子进程处理字节时更新计数器,主进程启动后台线程定期读取计数器值并更新进度条,实现平滑显示。
修改后的代码示例:
from concurrent.futures import ProcessPoolExecutor, as_completed from pathlib import Path from tqdm import tqdm import os import time import multiprocessing from threading import Thread class FileProcessor: NJOBS = os.get_cpu_count() CHK_SIZE = 1024 * 1024 * 80 # 80 MB def __init__(self, fp: str | Path): self.filepath = Path(fp) self.filesize = os.path.getsize(self.filepath) # 创建跨进程共享的无符号长整型计数器,初始值0 self.progress_counter = multiprocessing.Value('L', 0) self.progressbar = tqdm(total=self.filesize) @classmethod def _process_chunk(cls, chunk: bytes, progress_counter: multiprocessing.Value, *somemoreargs): array = bytearray(chunk) for i in range(len(array)): # 模拟字节处理逻辑 time.sleep(0.0001) # 更新共享计数器,加锁保证原子操作 with progress_counter.get_lock(): progress_counter.value += 1 return bytes(array) def _update_progress(self): """后台线程定期同步进度到tqdm""" last_val = 0 while True: with self.progress_counter.get_lock(): current_val = self.progress_counter.value # 进度完成时退出循环 if current_val >= self.filesize: self.progressbar.update(self.filesize - last_val) break # 只更新增量部分 if current_val > last_val: self.progressbar.update(current_val - last_val) last_val = current_val time.sleep(0.1) # 控制更新频率,避免资源浪费 def perform(self): def subchunk(chunk: bytes): schk_size = len(chunk) // FileProcessor.NJOBS if not schk_size: schk_size = len(chunk) i = 0 while (schunk := chunk[i:i + schk_size]): yield schunk i += schk_size # 启动后台进度更新线程 progress_thread = Thread(target=self._update_progress, daemon=True) progress_thread.start() file = self.filepath.open(mode="rb") executor = ProcessPoolExecutor(max_workers=FileProcessor.NJOBS) with file, executor: while (chunk := file.read(FileProcessor.CHK_SIZE)): futures = [ executor.submit( FileProcessor._process_chunk, sc, self.progress_counter ) for sc in subchunk(chunk) ] for future in as_completed(futures): # 处理子进程返回结果 res = future.result() # 等待进度条更新完成 progress_thread.join() # 最终处理逻辑
方案二:使用tqdm原生多进程支持
tqdm提供了contrib.concurrent.process_map工具,可直接在多进程任务中集成进度条,无需手动管理进程间通信,代码更简洁:
from tqdm.contrib.concurrent import process_map from pathlib import Path import os import time class FileProcessor: NJOBS = os.get_cpu_count() CHK_SIZE = 1024 * 1024 * 80 # 80 MB def __init__(self, fp: str | Path): self.filepath = Path(fp) self.filesize = os.path.getsize(self.filepath) # 预生成所有子任务的chunk列表 self.chunks = [] with self.filepath.open(mode="rb") as f: while (chunk := f.read(self.CHK_SIZE)): schk_size = len(chunk) // self.NJOBS if not schk_size: schk_size = len(chunk) i = 0 while (schunk := chunk[i:i + schk_size]): self.chunks.append(schunk) i += schk_size @classmethod def _process_chunk(cls, chunk: bytes): array = bytearray(chunk) for i in range(len(array)): # 模拟字节处理逻辑 time.sleep(0.0001) return bytes(array) def perform(self): # 使用process_map自动处理多进程和进度条 results = process_map( FileProcessor._process_chunk, self.chunks, max_workers=self.NJOBS, total=self.filesize, unit='B', unit_scale=True ) # 批量处理所有子任务结果 for res in results: # do something with res pass
方案说明
- 方案一适合需要自定义进度逻辑的场景,通过共享计数器保证进度准确性,后台线程负责平滑刷新进度条。
- 方案二利用tqdm原生工具,代码简洁易维护,适合任务结构清晰的场景。
注意事项:
- 共享计数器更新时必须加锁,避免多进程同时修改导致数据不一致。
- 后台线程设置为
daemon=True,确保主进程退出时线程自动终止。 - 方案二中
process_map的total参数需设置为文件总大小,保证进度条范围正确。
内容的提问来源于stack exchange,提问作者Michael Archman
相关产品推荐
相关产品推荐

