Apache Beam Python RunInference 变换实战:在 PCollection 上执行本地与远程机器学习推理

发布时间:2026/10/12 4:51:45
Apache Beam Python RunInference 变换实战:在 PCollection 上执行本地与远程机器学习推理 批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载RunInference 是 Apache Beam 面向机器学习推理场景的核心变换transform它直接作用于PCollection对批量或流式数据执行模型推理并输出输入样本 预测结果配对的结果集合。本文以 Beam 官方文档 runinference.md 为主体结合仓库源码与示例系统讲解 RunInference 的用法、PyTorch / Sklearn 两种框架下的完整代码、关键参数及底层实现原理。读完本文你将能够在自己的 Beam 流水线中快速接入本地或远程模型推理并理解 batching、指标收集、模型共享等进阶机制。什么是 RunInference 变换RunInference是apache_beam.ml.inference.base模块提供的变换对应 Python SDK 中的apache_beam.ml.inference.base.RunInference其核心职责是使用机器学习ML模型对PCollection中的一批样本examples执行推理inference并输出一个新的PCollection其中每个元素同时包含输入样本与模型预测结果。它同时支持本地推理模型加载到当前进程如 PyTorch、Sklearn、TensorFlow、ONNX、XGBoost 等与远程推理调用远端服务如 Vertex AI。根据官方文档的说明该 API 自Apache Beam 2.40.0 及后续版本可用。RunInference是一个泛型变换源码中其类型签名为见 base.pyclass RunInference(beam.PTransform[ beam.PCollection[Union[ExampleT, Iterable[ExampleT]]], beam.PCollection[PredictionT]]):也就是说输入是PCollection[样本]输出是PCollection[预测结果]样本与预测的具体类型由所选用的框架 ModelHandler 决定。两个核心抽象ModelHandler 与 PredictionResultModelHandler框架无关的模型加载与推理接口是RunInference的必填参数。它负责load_model()加载并初始化模型与run_inference()对一批样本执行推理两大能力。Beam 内置了针对 PyTorch、Sklearn、TensorFlow、ONNX、XGBoost、Vertex AI 等框架的 ModelHandler 实现见 sdks/python/apache_beam/ml/inference/ 目录。PredictionResult定义在 base.py 中的 NamedTuple包含三个字段example输入样本inference模型对该样本的预测结果model_id执行预测所用模型的标识通常是模型文件路径或 URI可选。从源码可以看出RunInference变换本身负责标准推理功能指标收集、在线程间共享模型、元素分批batching等见 base.py 的模块 docstring。官方示例总览Beam 官方文档为 RunInference 提供了两组框架示例每个框架都覆盖无键unkeyed模型与有键keyed模型两种数据形态框架示例PyTorchPyTorch 无键模型示例PyTorchPyTorch 有键模型示例SklearnSklearn 无键模型示例SklearnSklearn 有键模型示例这些示例使用一个公开的五倍表five times table线性模型输入x输出近似5 * x的预测值。示例代码位于仓库的sdks/python/apache_beam/examples/snippets/transforms/elementwise/目录下。示例一PyTorch 无键模型unkeyed无键模型指PCollection中每个元素就是一条独立的样本如一个 numpy 数组或 Tensor不需要附带键信息。完整示例见 runinference.pyimport apache_beam as beam import numpy import torch from apache_beam.ml.inference.base import RunInference from apache_beam.ml.inference.pytorch_inference import PytorchModelHandlerTensor model_state_dict_path gs://apache-beam-samples/run_inference/five_times_table_torch.pt model_class LinearRegression model_params {input_dim: 1, output_dim: 1} model_handler PytorchModelHandlerTensor( model_classmodel_class, model_paramsmodel_params, state_dict_pathmodel_state_dict_path) unkeyed_data numpy.array([10, 40, 60, 90], dtypenumpy.float32).reshape(-1, 1) with beam.Pipeline() as p: predictions ( p | InputData beam.Create(unkeyed_data) | ConvertNumpyToTensor beam.Map(torch.Tensor) | PytorchRunInference RunInference(model_handlermodel_handler) | beam.Map(print))要点拆解定义模型结构与参数示例中的LinearRegression是一个torch.nn.Module子类构造函数接收input_dim与output_dim通过model_params字典传入class LinearRegression(torch.nn.Module): def __init__(self, input_dim1, output_dim1): super().__init__() self.linear torch.nn.Linear(input_dim, output_dim) def forward(self, x): out self.linear(x) return out构造 ModelHandlerPytorchModelHandlerTensor同时接收model_class模型类、model_params实例化参数与state_dict_path模型权重文件路径支持 GCS、本地等 Beam FileSystems 可访问的 URI。state_dict_path与model_class必须成对出现二者缺一不可这一点由源码 pytorch_inference.py 的_validate_constructor_args校验若只传其一会抛出RuntimeError。数据转换模型输入类型是torch.Tensor因此先用beam.Map(torch.Tensor)将 numpy 数组转换为 Tensor再交给RunInference。接入流水线RunInference(model_handlermodel_handler)作为一个 PTransform 插入到流水线中输出被beam.Map(print)打印。运行后输出与测试用例 runinference_test.py 中check_torch_unkeyed_model_handler的期望一致PredictionResult(exampletensor([10.]), inferencetensor([52.2325]), model_idgs://apache-beam-samples/run_inference/five_times_table_torch.pt) PredictionResult(exampletensor([40.]), inferencetensor([201.1165]), model_idgs://apache-beam-samples/run_inference/five_times_table_torch.pt) PredictionResult(exampletensor([60.]), inferencetensor([300.3724]), model_idgs://apache-beam-samples/run_inference/five_times_table_torch.pt) PredictionResult(exampletensor([90.]), inferencetensor([449.2563]), model_idgs://apache-beam-samples/run_inference/five_times_table_torch.pt)注意模型是近似训练的五倍表模型因此预测值并非精确的5 * x而是带有误差的近似值。示例二PyTorch 有键模型keyed有键模型要求PCollection中的元素是(key, example)二元组推理完成后输出(key, prediction)便于将预测结果与原始输入一一对应。实现方式是使用KeyedModelHandler包装无键的 ModelHandler完整示例见 runinference.pyimport apache_beam as beam import torch from apache_beam.ml.inference.base import KeyedModelHandler from apache_beam.ml.inference.base import RunInference from apache_beam.ml.inference.pytorch_inference import PytorchModelHandlerTensor model_state_dict_path gs://apache-beam-samples/run_inference/five_times_table_torch.pt model_class LinearRegression model_params {input_dim: 1, output_dim: 1} keyed_model_handler KeyedModelHandler( PytorchModelHandlerTensor( model_classmodel_class, model_paramsmodel_params, state_dict_pathmodel_state_dict_path)) keyed_data [(first_question, 105.00), (second_question, 108.00), (third_question, 1000.00), (fourth_question, 1013.00)] with beam.Pipeline() as p: predictions ( p | KeyedInputData beam.Create(keyed_data) | ConvertIntToTensor beam.Map(lambda x: (x[0], torch.Tensor([x[1]]))) | PytorchRunInference RunInference(model_handlerkeyed_model_handler) | beam.Map(print))要点拆解输入数据是键值对列表键如first_question值是数值。用beam.Map(lambda x: (x[0], torch.Tensor([x[1]])))将值的部分转换为 Tensor保持键不变。KeyedModelHandler会把PCollection[Tuple[K, E]]转换为PCollection[Tuple[K, P]]源码见 base.py其内部先拆出键与样本推理后再把键与预测结果重新 zip 起来见 base.py。期望输出见 runinference_test.py(first_question, PredictionResult(exampletensor([105.]), inferencetensor([523.6982]), model_idgs://apache-beam-samples/run_inference/five_times_table_torch.pt)) (second_question, PredictionResult(exampletensor([108.]), inferencetensor([538.5867]), model_idgs://apache-beam-samples/run_inference/five_times_table_torch.pt)) (third_question, PredictionResult(exampletensor([1000.]), inferencetensor([4965.4019]), model_idgs://apache-beam-samples/run_inference/five_times_table_torch.pt)) (fourth_question, PredictionResult(exampletensor([1013.]), inferencetensor([5029.9180]), model_idgs://apache-beam-samples/run_inference/five_times_table_torch.pt))KeyedModelHandler还支持更高级的用法传入一组KeyModelMapping(keys, mh)让不同键的样本路由到不同的模型以及通过max_models_per_worker_hint限制每个 worker 进程同时驻留的模型数量避免多模型同时加载引发内存溢出OOM。这些能力同样定义在 base.py 的KeyedModelHandler中。示例三Sklearn 无键模型Sklearn 场景使用SklearnModelHandlerNumpy以 numpy 数组为输入并需要指明模型的序列化方式pickle 或 joblib。完整示例见 runinference_sklearn_unkeyed_model_handler.pyimport apache_beam as beam import numpy from apache_beam.ml.inference.base import RunInference from apache_beam.ml.inference.sklearn_inference import ModelFileType from apache_beam.ml.inference.sklearn_inference import SklearnModelHandlerNumpy sklearn_model_filename gs://apache-beam-samples/run_inference/five_times_table_sklearn.pkl sklearn_model_handler SklearnModelHandlerNumpy( model_urisklearn_model_filename, model_file_typeModelFileType.PICKLE) unkeyed_data numpy.array([20, 40, 60, 90], dtypenumpy.float32).reshape(-1, 1) with beam.Pipeline() as p: predictions ( p | ReadInputs beam.Create(unkeyed_data) | RunInferenceSklearn RunInference(model_handlersklearn_model_handler) | beam.Map(print))要点拆解model_uri指定模型文件位置model_file_type指定反序列化方式可选ModelFileType.PICKLE或ModelFileType.JOBLIB枚举定义见 sklearn_inference.py。若选择 JOBLIB 但运行环境未安装 joblib加载时会抛出ImportError见 sklearn_inference.py。Sklearn 场景不需要预先转换数据类型输入 numpy 数组直接进入RunInference。默认的 numpy 推理函数会用numpy.stack(batch, axis0)向量化一批样本后调用model.predict()见 sklearn_inference.py。期望输出见 runinference_test.pyPredictionResult(examplearray([20.], dtypefloat32), inferencearray([100.], dtypefloat32), model_idgs://apache-beam-samples/run_inference/five_times_table_sklearn.pkl) PredictionResult(examplearray([40.], dtypefloat32), inferencearray([200.], dtypefloat32), model_idgs://apache-beam-samples/run_inference/five_times_table_sklearn.pkl) PredictionResult(examplearray([60.], dtypefloat32), inferencearray([300.], dtypefloat32), model_idgs://apache-beam-samples/run_inference/five_times_table_sklearn.pkl) PredictionResult(examplearray([90.], dtypefloat32), inferencearray([450.], dtypefloat32), model_idgs://apache-beam-samples/run_inference/five_times_table_sklearn.pkl)示例四Sklearn 有键模型与 PyTorch 有键模型一样用KeyedModelHandler包装 Sklearn 的 ModelHandler。完整示例见 runinference_sklearn_keyed_model_handler.pyimport apache_beam as beam from apache_beam.ml.inference.base import KeyedModelHandler from apache_beam.ml.inference.base import RunInference from apache_beam.ml.inference.sklearn_inference import ModelFileType from apache_beam.ml.inference.sklearn_inference import SklearnModelHandlerNumpy sklearn_model_filename gs://apache-beam-samples/run_inference/five_times_table_sklearn.pkl sklearn_model_handler KeyedModelHandler( SklearnModelHandlerNumpy( model_urisklearn_model_filename, model_file_typeModelFileType.PICKLE)) keyed_data [(first_question, 105.00), (second_question, 108.00), (third_question, 1000.00), (fourth_question, 1013.00)] with beam.Pipeline() as p: predictions ( p | ReadInputs beam.Create(keyed_data) | ConvertDataToList beam.Map(lambda x: (x[0], [x[1]])) | RunInferenceSklearn RunInference(model_handlersklearn_model_handler) | beam.Map(print))与 PyTorch 版本唯一的差异在数据转换步骤Sklearn 期望每个样本是[单值列表]形式的特征向量因此使用beam.Map(lambda x: (x[0], [x[1]]))。期望输出见 runinference_test.py(first_question, PredictionResult(example[105.0], inferencearray([525.]), model_idgs://apache-beam-samples/run_inference/five_times_table_sklearn.pkl)) (second_question, PredictionResult(example[108.0], inferencearray([540.]), model_idgs://apache-beam-samples/run_inference/five_times_table_sklearn.pkl)) (third_question, PredictionResult(example[1000.0], inferencearray([5000.]), model_idgs://apache-beam-samples/run_inference/five_times_table_sklearn.pkl)) (fourth_question, PredictionResult(example[1013.0], inferencearray([5065.]), model_idgs://apache-beam-samples/run_inference/five_times_table_sklearn.pkl))从源码看 RunInference 的执行流程RunInference.expand()方法见 base.py揭示了变换内部的标准流水线结构大致分为五步前置处理preprocess依次执行 ModelHandler 注册的预处理函数如通过with_preprocess_fn添加将原始输入映射为底层 ModelHandler 期望的输入类型对应 PTransform 名称BeamML_RunInference_Preprocess。分批batching默认通过beam.BatchElements(**self._model_handler.batch_elements_kwargs())将元素聚合成批以提高模型调用效率。各框架的 ModelHandler 都支持min_batch_size、max_batch_size、max_batch_duration_secs三个分批参数见 pytorch_inference.py 与 sklearn_inference.py 中的_batching_kwargs构造逻辑若通过with_no_batching()关闭分批则输入必须已经是预分批的Iterable。执行推理核心是_RunInferenceDoFn见 base.py它在setup()阶段调用ModelHandler.load_model()加载模型并在process()中对每个批次调用ModelHandler.run_inference(batch, model, inference_args)得到预测结果。默认情况下模型在 DoFn 实例内通过shared.Shared复用避免每个 bundle 重复加载。后置处理postprocess依次执行 ModelHandler 注册的后处理函数如通过with_postprocess_fn添加对应 PTransform 名称BeamML_RunInference_Postprocess。结果装配run_inference()返回的预测结果通过utils._convert_to_result()见 utils.py逐条与输入样本 zip 成PredictionResult(example, inference, model_id)其中model_id默认取自模型文件路径或 URI。此外RunInference还提供几个实用构造参数见 base.pyinference_args传递给模型推理调用的额外参数只对需要额外参数的框架生效多数框架要求其为None。metrics_namespace指标收集的命名空间。model_metadata_pcoll/watch_model_pattern配合侧输入实现自动模型刷新Automatic Model Refresh。model_identifier为模型指定标识可在多个 RunInference 步骤间复用同一模型、避免重复加载需确保确实是同一个模型否则结果不确定。内置指标与模型共享机制RunInference默认通过_MetricsCollector见 base.py在指定命名空间下收集如下指标指标名类型含义num_inferencesCounter推理次数元素数failed_batches_counterCounter推理失败的批次数量inference_request_batch_sizeDistribution推理请求的批次大小inference_request_batch_byte_sizeDistribution推理请求的批次字节数inference_batch_latency_micro_secsDistribution批次推理延迟微秒model_byte_sizeDistribution模型字节数加载时估算load_model_latency_milli_secsDistribution模型加载延迟毫秒PyTorch 与 Sklearn 的 ModelHandler 分别使用BeamML_PyTorch与BeamML_Sklearn作为指标命名空间见 pytorch_inference.py 与 sklearn_inference.py。对于大模型可通过large_modelTrue或显式指定model_copies开启跨进程模型共享share_model_across_processes此时模型通过multi_process_shared.MultiProcessShared在多个进程中共享避免在同一台机器上加载多份大模型副本导致内存压力见 base.py 与 pytorch_inference.py。运行与验证以上四个示例均可在本地 Beam 环境直接运行也可以通过仓库中的单元测试验证行为。测试文件 runinference_test.py 使用TestPipeline和 mock 的print捕获输出逐行断言期望结果例如mock.patch(apache_beam.Pipeline, TestPipeline) mock.patch(apache_beam.examples.snippets.transforms.elementwise.runinference_sklearn_unkeyed_model_handler.print, str) class RunInferenceTest(unittest.TestCase): def test_sklearn_unkeyed_model_handler(self): runinference_sklearn_unkeyed_model_handler.sklearn_unkeyed_model_handler( check_sklearn_unkeyed_model_handler)需要说明的运行前提PyTorch 示例要求安装 torch测试代码通过try: import torch检测缺失时直接跳过见 runinference_test.py模型文件存放在 GCSgs://apache-beam-samples/...因此运行需要配置 GCP 文件系统依赖apache_beam.io.gcp.gcsfilesystem缺失时同样跳过测试。仓库中还提供了更多贴近真实业务的推理示例位于 sdks/python/apache_beam/examples/inference/ 目录包括 PyTorch 图像分类/分割、Sklearn 手写数字识别、TensorFlow、ONNX、XGBoost、Vertex AI 等场景可作为进一步学习的参考。小结RunInference 是 Apache Beam 将统一批流处理模型与机器学习推理结合的关键入口你只需要实现或选用一个ModelHandler即可把任意框架的模型推理接入既有 Beam 流水线并自动获得分批优化、线程间模型复用、指标监控、键值关联输出等能力。无论是本地加载的 PyTorch / Sklearn 模型还是后续通过侧输入与watch_model_pattern实现的自动模型刷新、远程 Vertex AI 服务其接入方式都遵循本文所述的同一套RunInference(model_handler...)模式。相关链接RunInference 官方文档RunInference PyTorch 示例文档RunInference Sklearn 示例文档RunInference 核心实现PyTorch ModelHandler 实现Sklearn ModelHandler 实现赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐Apache Beam Python RunInference 变换在批流统一管道中接入机器学习推理Apache Beam Python RunInference 变换在批流统一管道中接入机器学习推理 RunInference 是 Apache Beam 提Apache Beam RunInference API 实战在批流管道中运行机器学习模型推理Apache Beam RunInference API 实战在批流管道中运行机器学习模型推理 Apache Beam 通过 RunInference API批处理流处理大数据Plane 开源项目管理快速上手三步自托管1天跑通真实团队任务Plane 开源项目管理快速上手三步自托管1天跑通真实团队任务 Plane 是替代 Jira 和 Linear 的开源项目管理平台用来跟踪任务、跑迭代周期项目管理后端前端研发协作需求管理上一篇.NET 集合与 LINQ 性能优化实战指南从 FrozenDictionary 到零分配模式下一篇Uncloud 项目架构与开发指南读懂 AGENTS.md 背后的代码体系创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询