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

如何计算代码子进程运行耗时 实现Spacy批量文本处理ETA预估

Spacy 大规模文本处理的进度统计与ETA计算方案

核心问题说明

直接给nlp.pipe套通用进度条工具(比如alive-progress)效果差,核心原因有两个:

  • Spacy的nlp.pipe是批量处理逻辑,首次初始化模型、启动多子进程的耗时远高于常规单行迭代,通用进度条默认按匀速迭代估算ETA,偏差极大
  • 多进程模式下,通用进度条运行在主进程,没法直接捕获子进程内部的实际运行耗时,统计结果会混入进程调度、通信的冗余开销

具体实现步骤

1. 前置准备:内存友好的文件读取与模型加载

超大规模文本不要用readlines()一次性把全量内容读入内存,先逐行统计总行数作为进度计算基数,再用生成器逐行喂给处理逻辑:

from time import perf_counter
from spacy import load

# 配置项
MODEL_NAME = "zh_core_web_sm"  # 替换为你实际使用的Spacy模型
FILE_PATH = "text.txt"
PROCESS_NUM = 4  # 多进程数量,建议和CPU核心数匹配
LOG_INTERVAL = 2000  # 每处理多少行输出一次进度

# 逐行统计文件总行数,内存占用恒定
total_lines = 0
with open(FILE_PATH, "r", encoding="utf-8") as f:
    for _ in f:
        total_lines += 1

# 注意:多进程模式下不要在主进程加载模型,避免序列化开销

2. 轻量进度追踪包装器(无需第三方进度条依赖)

不要依赖第三方进度条的默认封装,自己实现迭代包装逻辑,基于实时处理速度动态计算ETA,统计结果更贴合Spacy的批量处理节奏:

def track_spacy_pipe(pipe_iter, total_count):
    """包装nlp.pipe迭代器,实时统计进度、速度、ETA"""
    start_ts = perf_counter()
    processed = 0

    # 时间格式化工具
    def _fmt_sec(sec):
        if sec < 60:
            return f"{sec:.1f}s"
        if sec < 3600:
            return f"{int(sec//60)}m{int(sec%60)}s"
        return f"{int(sec//3600)}h{int((sec%3600)//60)}m"

    for doc in pipe_iter:
        processed += 1
        # 到达日志间隔时输出统计信息
        if processed % LOG_INTERVAL == 0:
            cur_ts = perf_counter()
            total_cost = cur_ts - start_ts
            process_speed = processed / total_cost
            remain_count = total_count - processed
            eta = remain_count / process_speed if process_speed > 0 else 999999

            # 原地刷新输出进度
            print(
                f"进度: {processed}/{total_count} ({processed/total_count*100:.1f}%) | "
                f"已耗时: {_fmt_sec(total_cost)} | "
                f"速度: {process_speed:.0f}行/s | "
                f"ETA: {_fmt_sec(eta)}",
                end="\r",
                flush=True
            )
        yield doc

    # 处理完成输出最终统计
    final_cost = perf_counter() - start_ts
    print(f"\n全部处理完成!总耗时{_fmt_sec(final_cost)},平均处理速度{total_lines/final_cost:.0f}行/s")

# 调用示例
if __name__ == "__main__":
    nlp = load(MODEL_NAME)
    with open(FILE_PATH, "r", encoding="utf-8") as f:
        # 用生成器逐行喂给nlp.pipe,避免全量加载
        line_gen = (line.strip() for line in f if line.strip())
        docs = list(track_spacy_pipe(nlp.pipe(line_gen, n_process=PROCESS_NUM), total_lines))

3. 子进程耗时精准统计方案

如果需要单独统计每个子进程的实际运行耗时,不要用Spacy自带的n_process参数,手动实现进程池逻辑,在子进程内部打点计时,避免把主进程调度、IPC通信的时间算入子进程耗时:

from multiprocessing import Pool, current_process
from collections import defaultdict

# 子进程初始化函数:每个进程启动时只加载一次模型,避免重复开销
def _worker_init(model_name):
    global worker_nlp
    worker_nlp = load(model_name)
    global worker_start_ts
    worker_start_ts = perf_counter()  # 子进程内部计时起点
    print(f"子进程{current_process().pid}启动完成")

# 单行文本处理逻辑
def _process_single_line(line):
    pid = current_process().pid
    # 执行tokenize等NLP处理
    doc = worker_nlp(line.strip())
    # 统计当前子进程累计运行时长
    worker_runtime = perf_counter() - worker_start_ts
    return doc, pid, worker_runtime

# 主进程调用逻辑
if __name__ == "__main__":
    worker_costs = defaultdict(float)
    docs = []
    start_ts = perf_counter()
    processed = 0

    def _fmt_sec(sec):
        if sec < 60:
            return f"{sec:.1f}s"
        if sec < 3600:
            return f"{int(sec//60)}m{int(sec%60)}s"
        return f"{int(sec//3600)}h{int((sec%3600)//60)}m"

    with open(FILE_PATH, "r", encoding="utf-8") as f:
        line_gen = (line for line in f if line.strip())
        with Pool(processes=PROCESS_NUM, initializer=_worker_init, initargs=(MODEL_NAME,)) as pool:
            for doc, pid, runtime in pool.imap(_process_single_line, line_gen, chunksize=100):
                docs.append(doc)
                worker_costs[pid] = runtime
                processed +=1
                if processed % LOG_INTERVAL ==0:
                    total_cost = perf_counter() - start_ts
                    speed = processed/total_cost
                    eta = (total_lines - processed)/speed if speed>0 else 999999
                    print(
                        f"进度: {processed}/{total_lines} ({processed/total_lines*100:.1f}%) | "
                        f"速度: {speed:.0f}行/s | ETA: {_fmt_sec(eta)}",
                        end="\r", flush=True
                    )
    
    print(f"\n处理完成,总耗时{_fmt_sec(perf_counter()-start_ts)}")
    print("各子进程运行耗时统计:")
    for pid, cost in worker_costs.items():
        print(f"子进程{pid}: 累计运行{_fmt_sec(cost)}")

避坑提示

  • 不要一次性读取全量文件到内存,逐行生成器的内存占用不会随文件大小上涨,适合TB级文本处理场景
  • 多进程模式下必须在子进程初始化逻辑中加载模型,主进程加载模型后传入子进程会触发全量模型序列化传输,速度会比单进程还慢
  • ETA计算不要用初始批次的速度做固定值,Spacy首批处理包含模型初始化、子进程启动开销,速度会远低于稳定处理阶段,必须用滚动实时速度动态计算ETA
  • 子进程耗时不要在主进程侧计时,必须在子进程内部打点,否则统计结果会包含大量非处理逻辑的冗余开销,参考价值极低

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 04:57:10