如何计算代码子进程运行耗时 实现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
相关产品推荐
相关产品推荐

