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

如何结合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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 10:52:54