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

多进程批量处理数据时tqdm进度条修复及相关技术咨询

问题解答

1. 这种并行处理的批量实现方式是否正确且高效?

你这段代码里有个致命问题:vals = list(itertools.product(a,b))直接把所有元素转成列表,完全违背了“超大数组装不下内存”的初衷——如果a、b足够大,这一步直接就会内存溢出,所以这部分是错误的。

再说到效率:你的foo函数只是简单的a+b,计算量极小,但多进程的核心开销是进程间的数据传递,每个元素都要从主进程传到子进程,这种小计算量任务的通信开销会远大于并行带来的收益,反而比单进程运行更慢。

正确的优化方向:

  • 不要把itertools.product转成列表,直接用迭代器分批处理,避免一次性加载所有数据到内存。比如把vals = list(...)改成vals = itertools.product(a,b),同时修改chunks函数适配迭代器(因为迭代器没有len属性)。
  • 对于这种计算量极小的任务,要么合并更大的chunk,减少进程间通信的次数;要么直接用numpy的向量化操作代替循环——比如a[:, None] + b直接生成所有组合的和,再求和,效率比多进程高得多。
  • 如果必须用多进程,建议用pool.imap时指定chunksize参数,让进程池自动处理分块,减少主进程的管理开销,而不是手动分chunk。

2. 如何修复tqdm进度条,使所有批次的进度都在同一条进度条上更新?

当前代码每个chunk都创建新的tqdm实例,还把每个tqdm的total设成总任务数,导致每个chunk都打印一条从0到100%的进度条,自然会占多行。

给你两种修复方案:

方案一:手动分chunk + 总进度条

先修改chunks函数适配迭代器:

def chunks(iterable, size):
    chunk = []
    for item in iterable:
        chunk.append(item)
        if len(chunk) == size:
            yield chunk
            chunk = []
    if chunk:
        yield chunk

然后主函数里创建一个总进度条,每次处理完chunk就更新进度:

def main():
    np.random.seed(1234)
    
    a = np.random.rand(100)
    b = np.random.rand(100)
    total_tasks = len(a) * len(b)
    size_of_chunk = 1000
    # 用迭代器,不转列表
    vals = itertools.product(a,b)
                
    with mtp.Pool(processes=4) as pool:
        s = 0
        # 创建总进度条
        with tqdm(total=total_tasks) as pbar:
            for chunk in chunks(vals, size_of_chunk):
                results = list(pool.imap(foo, chunk))
                s += np.sum(results)
                # 更新进度条,步长是当前chunk的长度
                pbar.update(len(chunk))
        return s

方案二:直接用tqdm包裹整个任务迭代器(更简洁)

不用手动分chunk,让tqdm直接跟踪pool.imap的总进度:

def main():
    np.random.seed(1234)
    
    a = np.random.rand(100)
    b = np.random.rand(100)
    total_tasks = len(a) * len(b)
    # 用迭代器,不转列表
    vals = itertools.product(a,b)
                
    with mtp.Pool(processes=4) as pool:
        # tqdm直接包裹pool.imap,指定总任务数
        results = list(tqdm(pool.imap(foo, vals), total=total_tasks))
        s = np.sum(results)
        return s

这种方式不会产生多行进度条,tqdm会在同一行实时更新总进度。

3. 有哪些工具可以检测并行运行的Python脚本的内存占用情况?

分三类给你说:

代码内实时检测

用psutil库,可以获取主进程和所有子进程的内存占用,适合在代码关键位置打印监控:

import psutil
import os

def get_total_memory_mb():
    # 获取当前主进程
    main_proc = psutil.Process(os.getpid())
    # 获取所有子进程(包括进程池的子进程)
    child_procs = main_proc.children(recursive=True)
    # 计算总内存(rss是实际占用的物理内存,单位字节)
    total_rss = main_proc.memory_info().rss + sum(p.memory_info().rss for p in child_procs)
    # 转成MB
    return total_rss / (1024 * 1024)

比如在主循环里每隔几个chunk调用一次print(get_total_memory_mb()),就能看到内存变化。

系统级工具

  • Linux/macOS:用htop(比默认的top更直观),能看到所有Python进程的内存、CPU占用;也可以用ps aux | grep python过滤出Python进程,查看%MEM列的内存占比。
  • Windows:打开任务管理器,切换到“详细信息”标签,找到所有Python.exe进程,查看“内存(专用工作集)”列的数值。

Python专用工具

  • memory_profiler:可以逐行分析内存占用,多进程下需要在子进程的函数上加上@profile装饰器,或者用mprof run your_script.py命令运行脚本,生成内存使用的时间线报告。
  • tracemalloc:Python内置模块,能跟踪内存分配情况,适合排查内存泄漏,多进程下每个子进程需要单独初始化tracemalloc.start()。
  • py-spy:一款无需修改代码的采样分析工具,支持查看进程的CPU和内存使用情况,运行py-spy top --pid <你的Python进程ID>就能实时监控,包括子进程的情况。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 13:30:04