
Anomalib Pipelines 参考指南Pipeline、Job、JobGenerator 与 Runner 四大基础组件及 Benchmarking 流水线实战【免费下载链接】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 官方文档中的 Pipelines 参考页docs/source/markdown/guides/reference/pipelines/index.md展开系统讲解 Anomalib 实验性流水线的核心术语与四大基础类Pipeline、Job、JobGenerator、Runner并结合仓库源码深入剖析 Benchmarking 流水线的 YAML 网格搜索配置、串行/并行调度机制与结果落盘逻辑。读完本文你可以独立使用anomalib benchmark子命令跑多模型 × 多数据集的基准评测并具备基于四大基础类编写自定义流水线的能力。注意原文档醒目提示Pipelines 功能处于**实验性experimental**阶段未来版本可能在不保证向后兼容的情况下发生变更。这一声明同样体现在 CLI 入口的实现中——src/anomalib/cli/pipelines.py 中run_pipeline每次执行都会打印 “This feature is experimental. It may change or be removed in the future.” 的警告。一、为什么需要 Pipelines多模型、多阶段任务的编排问题Anomalib 文档给出的动机非常直接Benchmarking基准评测、Ensemble Tiling分块集成、Hyper-parameter optimization超参数优化这类任务本质上都需要“运行多个模型”并“把多个阶段串联起来”。Pipelines 特性提供了一种定义并执行此类任务的方式流水线的每个部分都被设计为独立且可组合independent and composable从而可以跨不同流水线复用。从源码结构看这一理念被落实在src/anomalib/pipelines/的目录划分中components/与具体业务无关的通用组件base 基础类、runners、utilsbenchmark/第一个落地的具体流水线pipeline、job、generator 三件套tiled_ensemble/面向分块集成的训练/测试流水线train_pipeline.py、test_pipeline.py其 How-To 教程见 docs/source/markdown/guides/how_to/pipelines/index.mdtypes.py定义PREV_STAGE_RESULT、RUN_RESULTS、GATHERED_RESULTS等跨阶段传递结果的类型别名。二、核心术语Pipeline、Runner、Job Generator、Job官方文档用四个术语勾勒出流水线的协作关系这也是理解全部源码的钥匙术语职责对应基类Pipeline主实体定义要执行的 Job 序列负责创建并运行 JobJob 由 Job Generator 生成阶段之间通过 Runner 串联PipelineRunner负责 Job 的调度与执行传递上一阶段的输出如存在调用正确的钩子hooks收集并保存 Job 结果再把汇总结果传给下一个 RunnerRunnerJob Generator根据配置生成 Job被 Runner 用来创建具体任务实例JobGeneratorJob原子工作单元可独立运行负责单个任务如训练一个模型或计算一次指标。其设计目标是不修改 Job 本身即可挂到任意 Runner 上方便将 Job 分布到不同执行者以提升吞吐Job文档同时建议想创建自定义流水线时参考 How-To Guide。三、四大基础类详解附源码级行为3.1 Pipeline 基类配置解析与 Runner 编排Pipeline是抽象基类定义了三个关键方法接口约定见 pipeline.py_get_args(args)若未传入args会通过get_parser()创建的jsonargparse.ArgumentParser从命令行解析必填参数--config随后用yaml.safe_load读取配置文件返回 dict。_setup_runners(args)抽象方法子类必须实现根据配置返回list[Runner]——流水线的并行/串行策略在这里决定。run(argsNone)核心执行循环。从源码pipeline.py#L74-L95可以看到关键行为日志统一重定向到runs/pipeline.log模块级常量log_file runs/pipeline.log按顺序遍历 runners每个 runner 的作业参数取自配置中runner.generator.job_class.name对应的顶级键例如 benchmark 流水线中 Job 的name benchmark因此读取config[benchmark]下的所有键作为job_argsprevious_results runner.run(job_args or {}, previous_results)上一 Runner 的汇总结果作为prev_stage_results传给下一个 Runner实现阶段间数据传递单个 Runner 抛异常不会中断整个流水线而是记录到日志并提示 “Please check runs/pipeline.log for more details”。get_parser静态方法是自定义流水线的参数入口不传 parser 时会新建一个只声明--config的jsonargparse解析器。3.2 Job 基类run / collect / save 三段式接口Jobjob.py#L42-L77强制子类实现三个部分类属性name: strJob 的标识同时是配置文件中对应顶级键名Runner 依此取参run(self, task_id: int | None None)执行核心逻辑。文档明确指出task_id是可选的仅在 Job 被并行执行时传入——这正是“Job 可挂到任意 Runner 而不需要改动自身”的关键设计串行执行时task_id为None并行执行时由ParallelRunner分配进程号collect(results)静态方法汇总所有run的返回值合并多份结果或做个体处理save(results)静态方法把汇总结果写入文件或数据库。collect与save均由 Runner 统一调用串行、并行两种 Runner 都遵循这一钩子约定因此 Job 作者只需专注run的业务逻辑。3.3 JobGenerator 基类把配置解析成 Job 迭代器JobGeneratorjob.py#L80-L109的任务是“解析配置并返回特定 Job 的迭代器”generate_jobs(args, prev_stage_result)抽象方法基于参数返回 Job 迭代器可用于展开网格组合、读取上一阶段结果等__call__只是generate_jobs的转发使generator(args, prev_stage_result)可以直接在for循环中使用job_class属性抽象属性返回将生成的 Job 类型Runner 依赖它调用collect/save。3.4 Runner 基类与两个内置实现Runner基类runner.py#L40-L72构造函数只接收一个generator: JobGenerator并定义抽象方法run(args, prev_stage_resultsNone)。文档中的args示例值得注意若流水线配置为arg1: arg2: hpo: param1: param2:则对应 HPO 生成器的args只会拿到hpo键下的所有键值——即“配置文件的顶级键 各 Job 类型各自的参数命名空间”。Anomalib 内置两个 RunnerSerialRunnerserial.py#L48-L127用tqdm进度条逐个执行generator(args, prev_stage_results)产出的 Job单个 Job 失败不会立即中断记录异常、置位failures其余 Job 继续执行全部执行完统一调用job_class.collect(results)与job_class.save(gathered_result)若存在失败打印/记录错误并抛出SerialExecutionError同时返回汇总结果。ParallelRunnerparallel.py#L46-L128文档docs/source/markdown/guides/reference/pipelines/runners/parallel.md说明并行 Runner 创建一个大小等于n_jobs的进程池池中每个进程被分配0到n_jobs-1的进程号执行 Job 时该编号会传给 JobJob 可用它实现进程专属逻辑——典型用法是当池大小等于 GPU 数量时用task_id把 Job 固定到指定 GPU。源码印证了这一设计# src/anomalib/pipelines/components/runners/parallel.py with ProcessPoolExecutor(max_workersself.n_jobs, mp_contextmultiprocessing.get_context(spawn)) as executor: 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)要点使用spawn多进程上下文维护processes: dict[int, Future | None]作为空闲槽位表有空槽才提交新 Job_await_cleanup_processes负责回收已完成进程的结果单个进程异常只置位failures不拖垮整个池全部结束后同样走collect→save钩子若failures为真则抛出ParallelExecutionError。3.5 结果类型约定跨阶段传递的类型定义在 src/anomalib/pipelines/types.pyRUN_RESULTS单个 Job 的run返回值、GATHERED_RESULTScollect的返回值、PREV_STAGE_RESULT传给下一 Runner 的阶段结果在 base pipeline 中仅以TYPE_CHECKING方式引用。四、内置流水线之一Benchmarking参考文档列出的首个可用流水线是Benchmarking使用网格搜索grid-search跨模型计算指标入口页docs/source/markdown/guides/reference/pipelines/benchmark/index.md。4.1 配置结构与网格搜索一次基准运行由配置文件驱动该文件指定网格搜索参数。官方文档给出的示例配置accelerator: - cuda - cpu benchmark: seed: 42 model: class_path: grid_search: [Padim, Patchcore] data: class_path: MVTecAD init_args: category: grid: - bottle - cable - capsule配置遵循 jsonargparse 约定class_pathinit_args用于实例化对象MVTecAD数据模块、Padim/Patchcore模型均通过get_model/get_datamodule由类名构建见 generator.py#L99-L108grid键表示该参数要展开为多个值每个组合生成一个独立 Job。仓库中也提供了一份可直接运行的示例 tools/experimental/benchmarking/sample.yamlPadim/Patchcore×MVTec的bottle/capsule网格accelerator 同为cudacpu。网格展开的底层实现是 get_iterator_from_grid_dict先用flatten_dict把配置压平成点号键筛出键名含grid的项用itertools.product生成全部组合逐个回填后再to_nested_dict还原嵌套结构——非网格参数如seed原样保留。文档示例中 2 个模型 × 3 个类别即会展开为 6 个 Job。accelerator是流水线专属参数用于配置 Runner传入cuda时流水线会添加一个 Job 数等于 CUDA 设备数的 ParallelRunnercpu的 Job 则 串行执行。设计动机与文档表述一致Job 彼此独立把每个 Job 分发到独立加速器可以提升吞吐。4.2 Benchmark 流水线的 Runner 装配逻辑Benchmark._setup_runnersbenchmark/pipeline.py#L90-L125把上述策略代码化只支持cpu与cuda其他值抛ValueError: Unsupported accelerator调用torch.cuda.device_count()若设备数 1或加速器是cpu装配SerialRunner(BenchmarkJobGenerator(accelerator))否则装配ParallelRunner(BenchmarkJobGenerator(accelerator), n_jobsdevice_count)即每块 GPU 一个进程配置里每个 accelerator 各自装配一个 RunnerPipeline.run会按列表顺序依次执行——因此示例配置中cuda组先并行跑完cpu组再串行执行。4.3 BenchmarkJob单个 Job 做了什么BenchmarkJob 的name benchmark与配置顶级键对应。runjob.py#L111-L164的执行流程task_id不为None时并行模式把devices设为[task_id]实现“进程号 → 指定 GPU”的绑定在临时目录中seed_everything(self.seed)保证可复现构造Engine(accelerator..., devices..., default_root_dirtemp_dir)依次计时执行engine.fit(model, datamodule)与engine.test(model, datamodule)返回字典包含accelerator、三段计时job_duration/fit_duration/test_duration、展平后的配置flat_cfg以及test产出的全部指标。汇总与落盘collect把各 Job 结果字典转成pandas.DataFrame每行一次基准运行save先用rich.table.Table在终端打印结果表再把 CSV 写到runs/benchmark/YYYY-MM-DD-HH_MM_SS/results.csv目录不存在会自动创建。4.4 BenchmarkJobGenerator如何把配置变成 JobBenchmarkJobGenerator 持有构造时的acceleratorjob_class属性返回BenchmarkJob。generate_jobs对get_iterator_from_grid_dict(args)的每一个组合yield一个BenchmarkJob( acceleratorself.accelerator, seed_container[seed], modelget_model(_container[model]), datamoduleget_datamodule(_container[data]), flat_cfgflatten_dict(_container), )previous_stage_result在此未被使用del previous_stage_result说明 benchmark 是单阶段流水线但接口上仍保留了多阶段扩展位。4.5 运行 Benchmarking 流水线文档给出两种运行方式方式一Anomalib 子命令anomalib benchmark --config tools/experimental/benchmarking/sample.yaml方式二独立入口脚本python tools/experimental/benchmarking/benchmark.py --config tools/experimental/benchmarking/sample.yaml两种方式的差异在源码中清晰可见CLI 子命令由 src/anomalib/cli/pipelines.py 中的注册表PIPELINE_REGISTRY {benchmark: Benchmark}驱动run_pipeline取出子命令对应类并调用run(config)而tools/experimental/benchmarking/benchmark.py是直接import Benchmark后调用run的独立入口。另外Pipeline.run在无参调用时会自动从sys.argv解析--config见 benchmark/pipeline.py 模块 docstring 中的两种用法示例因此也可以在脚本中写Benchmark().run()直接消费命令行参数。执行完毕后除终端的 rich 表格外还会在runs/benchmark/时间戳/results.csv得到结构化结果、在runs/pipeline.log得到全量日志便于二次分析。五、相关文档与延伸阅读参考页自身的文档树指向以下同级页面可作为继续深入的入口均相对仓库根目录Pipeline 基类参考、Job 参考、Job Generator 参考Runner 参考含 Serial Runner 与 Parallel RunnerBenchmark 流水线含 Benchmark Job 与 Benchmark Job GeneratorHow-To 教程索引Tiled Ensemble 与自定义流水线Custom Pipeline的实操指南。六、小结设计要点速览配置即编排配置文件顶级键 Job 类型Job.name的参数空间grid键触发笛卡尔积展开每个组合一个 JobJob 与 Runner 解耦task_id可选参数让同一个 Job 无需修改即可被串行或并行调度这是吞吐扩展如按 GPU 数开进程的前提钩子统一收敛collect/save由 Runner 统一调用保证串行/并行两条路径的结果落盘行为一致benchmark 场景即runs/benchmark/时间戳/results.csv;实验性声明整套 Pipelines 机制包括anomalib benchmark子命令均为实验特性升级版本时需关注其接口变更。【免费下载链接】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),仅供参考