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 目录Jobsrc/anomalib/pipelines/components/base/job.py原子工作单元子类必须实现run、collect、save三个方法JobGenerator负责解析配置参数产出具体 Job 实例的迭代器并通过job_class属性声明其生成的 Job 类型Runnersrc/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_workersself.n_jobs, mp_contextmultiprocessing.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_idindex) self._await_cleanup_processes(blockingTrue) # 等待全部任务收尾可以看到调度策略是“有空位就填、没空位就等”_await_cleanup_processes(blockingFalse)轮询processes表一旦某Future完成即调用.result()取回结果、追加到self.results并把该槽位置回None所有 Job 提交完后再以blockingTrue等待收尾。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(fcuda:{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_jobstorch.cuda.device_count()) results runner.run({param: value})n_jobstorch.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 fThere 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_jobsdevice_count))规则清晰单卡或 CPU 时并行没有收益退回串行多卡时按卡数开池并行评测。2. Tiled Ensemble 管道对每个阶段做相同判断。train_pipeline.py#L79-L109 中训练阶段与验证预测阶段都遵循accelerator cuda则ParallelRunner(..., n_jobstorch.cuda.device_count())、否则SerialRunner的分支而“合并预测”“缝合平滑”这类天然是单进程后处理的阶段则固定使用SerialRunner。测试管道 test_pipeline.py 采用同样的模式。这展示了实际工程中“阶段级”的执行器选择粒度并行只用于可独立并行的批处理阶段。3. 自定义管道中通过配置声明并行。文档示例 docs/source/snippets/pipelines/dummy/pipeline_parallel.txt 展示了管道配置中直接内嵌ParallelRunner(TrainJobGenerator(), n_jobsargs[train][experiments])的写法——即池大小可以绑定到配置里的实验数量而不必是 GPU 数。该示例属于 how-to 指南 自定义管道 的一部分文中明确指出“由于所有任务相互独立可以使用 ParallelRunner”。并行 / 串行选择建议结合源码行为可以归纳如下实践准则场景推荐执行器依据任务相互独立、每任务独占一张 GPUParallelRunner(n_jobsdevice_count)任务内用task_id绑定设备内置 benchmark / tiled_ensemble 管道均采用此模式单卡或 CPUSerialRunner并行无加速收益且免去 spawn 进程开销任务有顺序依赖或需调试SerialRunner串行输出可预期且串行版带 tqdm 进度条见 serial.py#L113池大小不必等于 GPU 数如按实验数开池ParallelRunner(n_jobs配置值)任务内部自行决定设备策略见 dummy 管道示例片段需要记住的前提由于使用spawn上下文Job 实例及其args必须可被 pickletask_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),仅供参考
