Anomalib 管道并行执行:ParallelRunner 进程池机制与多 GPU 任务调度实践
【免费下载链接】anomalibAn anomaly detection library comprising state-of-the-art algorithms and features such as experiment management, hyper-parameter optimization, and edge inference.项目地址: https://gitcode.com/GitHub_Trending/an/anomalib
本篇指南聚焦 Anomalib 管道(pipeline)框架中的ParallelRunner:它通过一个大小可配置的多进程池并行执行管道任务(Job),并向每个任务注入 0 到n_jobs-1的进程 ID,使任务能够按 ID 独占指定 GPU 等资源。读完本文,你将理解ParallelRunner的进程池调度、结果汇聚与失败处理机制,并掌握它在基准测试(benchmark)与分块集成(tiled ensemble)等内置管道中的实际用法,以及如何在自定义管道 YAML 中配置并行执行。
ParallelRunner 在管道框架中的定位
Anomalib 的管道体系由三类抽象组件构成,代码位于 components 目录:
- Job(src/anomalib/pipelines/components/base/job.py):原子工作单元,子类必须实现
run、collect、save三个方法; - JobGenerator:负责解析配置参数,产出具体 Job 实例的迭代器,并通过
job_class属性声明其生成的 Job 类型; - Runner(src/anomalib/pipelines/components/base/runner.py):定义
run(args, prev_stage_results)抽象接口,决定 Job 的执行方式——串行或并行。
Runner.run的两个参数来自管道配置:args是当前阶段配置块下的全部键值(例如 HPO 阶段拿到hpo节点下的参数字典),prev_stage_results是上一阶段的汇聚结果,用于阶段间依赖。Runner 只负责“怎么跑”,Job 负责“跑什么”,两者解耦使同一组 Job 可以在串行/并行执行器间自由切换——这正是 runners 文档索引页 中并列提供 Serial Runner 与 Parallel Runner 的原因。
进程池机制:pool size、进程 ID 与 spawn 上下文
ParallelRunner的完整实现位于 src/anomalib/pipelines/components/runners/parallel.py。其构造签名只有一个关键参数:
def __init__(self, generator: JobGenerator, n_jobs: int) -> None: super().__init__(generator) self.n_jobs = n_jobs self.processes: dict[int, Future | None] = {} self.results: list[dict] = [] self.failures = False核心语义与文档 parallel.md 一致:
- 进程池大小等于创建时定义的
n_jobs。run方法中通过ProcessPoolExecutor(max_workers=self.n_jobs, mp_context=multiprocessing.get_context("spawn"))创建池(parallel.py#L92-L113)。注意这里显式使用spawn而非默认的fork上下文:spawn以全新解释器进程启动 worker,避免fork在 GPU 驱动、多线程库等场景下的已知隐患,代价是启动更慢、要求 Job 对象可被 pickle。 - 每个进程持有 0 到
n_jobs-1的进程 ID。run一开始便初始化self.processes = dict.fromkeys(range(self.n_jobs)),即一个以进程 ID 为键、以Future(或空闲标记None)为值的状态表。 - 任务提交时进程 ID 被透传给 Job。调度主循环的逻辑是:
for job in self.generator(args, prev_stage_results): while None not in self.processes.values(): self._await_cleanup_processes() # 池满时等待任一进程结束并回收 index = next(i for i, p in self.processes.items() if p is None) self.processes[index] = executor.submit(job.run, task_id=index) self._await_cleanup_processes(blocking=True) # 等待全部任务收尾可以看到调度策略是“有空位就填、没空位就等”:_await_cleanup_processes(blocking=False)轮询processes表,一旦某Future完成即调用.result()取回结果、追加到self.results并把该槽位置回None;所有 Job 提交完后再以blocking=True等待收尾。
task_id:让每个进程绑定独立 GPU
官方文档给出的核心用法:若进程池大小等于 GPU 数量,任务即可用进程 ID 直接指定所用 GPU。Job基类对这一机制的约定是——run(self, task_id: int | None = None)中task_id仅在并行执行时传入(串行执行不传),见 job.py#L47-L52。源码 docstring 中给出的 Job 侧示例:
def run(self, arg1: int, arg2: nn.Module, task_id: int) -> None: device = torch.device(f"cuda:{task_id}") # 进程 ID 即 GPU 序号 model = arg2.to(device) # ... 其余任务逻辑对应的 Runner 侧典型用法(来自模块 docstring 的官方示例):
from anomalib.pipelines.components.runners import ParallelRunner from anomalib.pipelines.components.base import JobGenerator import torch generator = JobGenerator() runner = ParallelRunner(generator, n_jobs=torch.cuda.device_count()) results = runner.run({"param": "value"})n_jobs=torch.cuda.device_count()使每个 worker 进程恰好独占一张卡,task_id与cuda:{task_id}一一对应,天然避免了多进程争抢同一 GPU 显存的问题。这个“池大小 = 设备数、ID = 设备号”的模式在仓库内置管道中被反复复用。
结果汇聚与失败处理
所有 worker 结束后,run方法并不自己合并结果,而是把职责交回给 Job 类的静态方法:
gathered_result = self.generator.job_class.collect(self.results) self.generator.job_class.save(gathered_result) if self.failures: msg = f"There were some errors with job {self.generator.job_class.name}" print(msg) logger.error(msg) raise ParallelExecutionError(msg)即 job.py 中约定的三段式接口:
collect(results: list[RUN_RESULTS]):把各进程run的返回值合并为一份结构化结果;save(gathered_result):落盘或入库;- 失败路径:任一进程
Future.result()抛异常时,_await_cleanup_processes会记录日志并置self.failures = True;注意它不会中断其余任务,而是等全部任务跑完后再统一抛出ParallelExecutionError(与SerialRunner抛出的SerialExecutionError语义对齐,两者都在 serial.py 与 parallel.py 中分别定义)。
collect/save的“批量化”特性也解释了为何并行与串行 Runner 可以互换:两者最终都产出同一形态的GATHERED_RESULTS,下游管道阶段无需感知执行方式。
仓库中的实际用法
1. 基准测试管道按加速设备自动选择执行器。src/anomalib/pipelines/benchmark/pipeline.py 中:
device_count = torch.cuda.device_count() if device_count <= 1 or accelerator == "cpu": runners.append(SerialRunner(BenchmarkJobGenerator(accelerator))) else: runners.append(ParallelRunner(BenchmarkJobGenerator(accelerator), n_jobs=device_count))规则清晰:单卡或 CPU 时并行没有收益,退回串行;多卡时按卡数开池并行评测。
2. Tiled Ensemble 管道对每个阶段做相同判断。train_pipeline.py#L79-L109 中,训练阶段与验证预测阶段都遵循accelerator == "cuda"则ParallelRunner(..., n_jobs=torch.cuda.device_count())、否则SerialRunner的分支,而“合并预测”“缝合平滑”这类天然是单进程后处理的阶段则固定使用SerialRunner。测试管道 test_pipeline.py 采用同样的模式。这展示了实际工程中“阶段级”的执行器选择粒度:并行只用于可独立并行的批处理阶段。
3. 自定义管道中通过配置声明并行。文档示例 docs/source/snippets/pipelines/dummy/pipeline_parallel.txt 展示了管道配置中直接内嵌ParallelRunner(TrainJobGenerator(), n_jobs=args["train"]["experiments"])的写法——即池大小可以绑定到配置里的实验数量,而不必是 GPU 数。该示例属于 how-to 指南 自定义管道 的一部分,文中明确指出“由于所有任务相互独立,可以使用 ParallelRunner”。
并行 / 串行选择建议
结合源码行为,可以归纳如下实践准则:
| 场景 | 推荐执行器 | 依据 |
|---|---|---|
| 任务相互独立、每任务独占一张 GPU | ParallelRunner(n_jobs=device_count),任务内用task_id绑定设备 | 内置 benchmark / tiled_ensemble 管道均采用此模式 |
| 单卡或 CPU | SerialRunner | 并行无加速收益,且免去 spawn 进程开销 |
| 任务有顺序依赖或需调试 | SerialRunner | 串行输出可预期,且串行版带 tqdm 进度条(见 serial.py#L113) |
| 池大小不必等于 GPU 数(如按实验数开池) | ParallelRunner(n_jobs=<配置值>),任务内部自行决定设备策略 | 见 dummy 管道示例片段 |
需要记住的前提:由于使用spawn上下文,Job 实例及其args必须可被 pickle;task_id只在并行执行时传入,Job 的run签名应保持task_id: int | None = None的可选形式,才能同时兼容串行执行器。
参考文档:Parallel Runner、Runner 索引、Serial Runner。
【免费下载链接】anomalibAn anomaly detection library comprising state-of-the-art algorithms and features such as experiment management, hyper-parameter optimization, and edge inference.项目地址: https://gitcode.com/GitHub_Trending/an/anomalib
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考