大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 在统一的批流编程模型中内置了完整的 AI/ML 能力既可以用同一套 PTransform 处理大规模数据集的预处理与模型推理也可以在 MLOps 生态中承担探索性数据分析EDA到生产流水线平滑扩容的职责。本文将围绕 Beam 官方文档中关于 AI/ML 能力的说明深入解析以RunInference API为核心的模型加载、推理、批量配置、自定义模型接入与流式自动模型刷新机制并结合本仓库 Python SDK 的源码实现帮助你掌握一套可落地的端到端推理方案。读完本文你将掌握如何在 Beam 管道中通过ModelHandler封装并加载 PyTorch / scikit-learn / TensorFlow 等框架的预训练模型如何使用RunInferencePTransform 对PCollection中的样本批量执行推理如何在流式管道中借助WatchFilePattern与 side input 实现不停机自动切换最新模型。Apache Beam 内置 AI/ML 能力概览根据 Beam 官方文档本文所依据的关联文档 17_advanced_ai_ml.mdApache Beam 提供了三类核心 AI/ML 能力大规模数据预处理与模型推理对大型数据集同时进行特征工程、清洗、聚合等预处理并在同一管道内完成模型推理避免数据在不同系统间搬运。探索性数据分析与 MLOps 生产化以 Beam 管道完成数据探索随后将同一套逻辑平滑扩展到生产环境作为 MLOps 生态中的数据处理与推理环节。批流一体的生产推理在数据负载变化的情况下于批处理与流处理管道中运行模型满足在线服务与离线批推的混合场景。官方文档明确给出两条推荐路径推理实现的首选方式是 RunInference API它统一封装了 PyTorch、scikit-learn、TensorFlow 等框架的推理逻辑对于常见的云平台集成模式如 Vertex AI 等可参考官方文档中的 AI Platform integration patterns 章节。Beam 同时支持使用 PyTorch、scikit-learn、TensorFlow 的预训练模型也支持自定义模型custom models的推理接入。RunInference API统一推理入口RunInference 是 Apache Beam 在apache_beam.ml.inference包中提供的核心 PTransform。它接收一个PCollection的样本examples使用 ML 模型执行推理输出一个同时包含输入样本与对应预测结果的PCollection[PredictionResult]。其 Python SDK 实现位于 base.py核心类定义为class RunInference(beam.PTransform[beam.PCollection[Union[ExampleT, Iterable[ExampleT]]], beam.PCollection[PredictionT]]): def __init__( self, model_handler: ModelHandler[ExampleT, PredictionT, Any], clocktime, inference_args: Optional[dict[str, Any]] None, metrics_namespace: Optional[str] None, *, model_metadata_pcoll: beam.PCollection[ModelMetadata] None, watch_model_pattern: Optional[str] None, model_identifier: Optional[str] None, use_model_manager: bool False, model_manager_args: Optional[dict[str, Any]] None, monitoring_transform: Optional[beam.PTransform] None, **kwargs):关键参数说明依据 base.py 源码model_handlerModelHandler的实现是必需参数负责加载模型并在内部对样本批量推理inference_args传递给模型推理调用的额外参数如部分框架推理函数所需的超参数metrics_namespace推理指标收集的命名空间默认由ModelHandler.get_metrics_namespace()提供默认值为RunInferencemodel_metadata_pcoll一个以 side input 形式传入的PCollection[ModelMetadata]用于自动刷新模型详见下文自动模型刷新一节watch_model_pattern一个 glob 模式用于监视目录并自动刷新模型与model_metadata_pcoll二选一的便捷写法model_identifier模型标识字符串可用于多个 RunInference 步骤间复用同一模型、避免重复加载需要注意不同模型使用同一 tag 会导致不确定结果use_model_manager/model_manager_args是否启用 ModelManager 进行多副本模型管理与资源伸缩。版本与 SDK 支持范围Python SDK 自2.40.0起提供 RunInference APIJava SDK 自2.41.0起通过 Apache Beam 的**多语言管道Multi-language Pipelines**框架支持该 API推理支持批处理与流处理两种模式并且支持 GPU 推理如 PyTorch 的 CUDA 环境变量可通过env_vars传入。支持的框架与模型中心RunInference 官方支持以下框架与模型中心依据文档 38_ai_runinference.md 及本仓库apache_beam/ml/inference目录下的实现文件框架 / 模型中心仓库实现文件PyTorchpytorch_inference.pyscikit-learnsklearn_inference.pyTensorFlowtensorflow_inference.pyXGBoostxgboost_inference.pyHugging Facehuggingface_inference.pyTensorFlow Hub通过 tensorflow_inference.py 的 hub 模块加载Vertex AIvertex_ai_inference.pyTensorRTtensorrt_inference.pyONNXonnx_inference.py此外本仓库还提供了 vLLMvllm_inference.py、Anthropicanthropic_inference.py、Geminigemini_inference.py等推理封装并在sdks/python/apache_beam/ml/inference目录下为每个框架配套了*_test.py单元测试与*_it_test.py集成测试可作为接入与验证的参考实现。ModelHandler模型加载与批量推理的抽象层ModelHandler是 RunInference 与具体模型框架之间的抽象层抽象基类定义于 base.py。一个ModelHandler实现需要提供两个核心能力load_model()加载并初始化模型run_inference(batch, model, inference_args)对一个样本批次执行推理并返回Iterable[PredictionT]。ModelHandler.__init__同时接受一组批量batching配置参数base.py 源码这些参数最终会被转换为beam.BatchElements的 kwargs参数含义min_batch_size/max_batch_size输入样本分批的最小 / 最大批量大小max_batch_duration_secs流式场景下缓冲一个批次的最长等待时间秒超时即发射max_batch_weight一个批次的最大权重需配合element_size_fn使用element_size_fn返回单个元素权重字节大小的函数batch_length_fn将元素映射为长度int的可调用对象用于变长输入的长度感知分桶减少 padding 浪费batch_bucket_boundaries长度分桶的有序边界值列表下界包含式bisect_right语义默认[16, 32, 64, 128, 256, 512]large_model模型大到多个副本会引发内存压力时置为True此时会跨进程共享模型model_copies期望在单机上加载的模型精确副本数env_varskwargs加载模型前设置的环境变量例如 GPU 推理所需的CUDA_VISIBLE_DEVICES从源码结构看RunInference会利用ModelHandler.batch_elements_kwargs()base.py取得上述批量参数并交给beam.BatchElements从而实现集中式模型管理、分块推理的内存/带宽优化——这正是 RunInference Centralized model management 特性的底层实现。以 PyTorch 为例接入预训练模型文档 38_ai_runinference.md 给出了导入 PyTorch 模型处理器的骨架代码from apache_beam.ml.inference.pytorch_inference import PytorchModelHandlerTensor from apache_beam.ml.inference.base import RunInference model_handler PytorchModelHandlerTensor( # Model handler setup model_pathgs://your-bucket/model.pth, # 模型路径或 URI model_classMyModelClass, # 自定义模型类 model_params{in_features: 128}, # 模型构造函数参数 env_vars{CUDA_VISIBLE_DEVICES: 0}, # 可选GPU 环境变量 ) with pipeline as p: predictions p | Read beam.ReadFromSource(a_source) | RunInference RunInference(model_handler)运行后RunInference输出的PCollection[PredictionResult]中每个元素是PredictionResult(example, inference, model_id)。PredictionResult在 base.py 中被定义为NamedTupleclass PredictionResult(NamedTuple(PredictionResult, [(example, _INPUT_TYPE), (inference, _OUTPUT_TYPE), (model_id, Optional[str])])): __slots__ () def __new__(cls, example, inference, model_idNone): return super().__new__(cls, example, inference, model_id)即输出同时携带原始输入样本与推理结果以及可选的model_id用于标识模型版本便于下游进行后处理、评估与指标统计。在管道中实现批流一体的模型推理RunInference 对批处理与流处理一视同仁样本以PCollection进入推理结果以PCollection输出上游可以是任何Read变换下游可以是任何Write、ParDo或聚合变换。流式场景下max_batch_duration_secs会控制批次在窗口内的缓冲时长从而在吞吐优先的大批次与延迟优先的小批次之间取得平衡。对于模型推理过程中的资源提示GPU/内存等ModelHandler.get_resource_hints()base.py可为变换附加资源提示供支持资源调度的 Runner如 Dataflow申请对应资源。自动模型刷新不停机更新生产模型生产 MLOps 场景中模型需要随新数据持续更新。Beam 官方文档41_ai_model_refresh.md给出的推荐做法是RunInference API Side Inputs让流式管道在不停止管道的情况下自动切换为最新模型。ModelMetadata 与 model_metadata_pcollRunInference接受可选参数model_metadata_pcoll这是一个以 side input 形式存在的PCollection其中每个元素是ModelMetadata。ModelMetadata定义于 base.pyclass ModelMetadata(NamedTuple): model_id: str model_name: strmodel_id模型唯一标识可以是文件路径或 URL用于加载模型执行推理model_name人类可读的模型名用于在 RunInference 生成的指标中标识模型。ModelMetadata.model_id与model_name的完整语义注释见 base.py。文档同时提醒model_id中的 URL/路径必须与对应ModelHandler的加载要求兼容若主输入先于 side input 到达主PCollection会被缓冲直到model_metadata_pcoll产出元素。使用 WatchFilePattern 监视模型目录文档给出的生产级更新模式是使用WatchFilePattern作为 side inputimport apache_beam as beam from apache_beam.ml.inference.utils import WatchFilePattern from apache_beam.ml.inference.base import RunInference tf_model_handler ... # model handler for the model with beam.Pipeline() as pipeline: file_pattern path_to_model_file side_input_pcoll ( pipeline | FilePatternUpdates WatchFilePattern(file_patternfile_pattern)) main_input_pcoll ... # main input PCollection inference_pcoll ( main_input_pcoll | RunInference RunInference( model_handlermodel_handler, model_metadata_pcollside_input_pcoll))WatchFilePattern的实现位于 utils.py其核心流程为MatchContinuously持续监视匹配file_pattern的文件支持本地路径与gs://GCS 路径含*、?、[...]glob 字符_GetLatestFileByTimeStamp依据文件修改时间挑选管道启动后更新过的最新文件并封装为ModelMetadata(model_idmodel_path, model_namemodel_name)_ConvertIterToSingleton借助状态state保证输出可被beam.pvalue.AsSingleton()包装为单值 side input最终WindowInto使用GlobalWindowsRepeatedly(AfterProcessingTime(1))DISCARDING累积模式将新模型信息周期性地广播给RunInference。WatchFilePattern的构造函数参数utils.pyfile_pattern模型文件的本地路径或gs://路径可含 glob 通配符interval检查文件匹配的间隔秒默认 360 秒stop_timestamp停止检查的时间戳默认MAX_TIMESTAMP即不停止。从源码注释可以提炼三条重要使用前提文件名不可复用若新增/更新文件使用了之前用过的文件名该变换会忽略这次更新触发模型更新必须上传唯一命名的新文件初始必须存在匹配文件管道启动前file_pattern至少要匹配到一个文件empty_match_treatmentEmptyMatchTreatment.DISALLOW仅适用于流式模式该变换依赖MatchContinuously产生无界数据源在批处理模式下可能产生非预期结果或导致管道卡住。另外文档强调model_metadata_pcoll参数期望的是一个与AsSingleton标记兼容的ModelMetadataPCollection而WatchFilePattern会自动管理窗口并把输出封装为ModelMetadata因此直接将其作为 side input 传入即可。接入自定义模型当目标模型不属于官方支持的框架时可以自行实现ModelHandler或KeyedModelHandler完成模型的加载与推理逻辑。文档提到的一个典型示例是使用 spaCy 加载自定义 NLP 模型。实现自定义ModelHandler只需继承抽象基类并实现两个方法base.pyclass MyModelHandler(ModelHandler[ExampleT, PredictionT, ModelT]): def load_model(self) - ModelT: # 1. 设置 env_vars # 2. 从 self._model_uri 加载模型并返回 ... def run_inference(self, batch, model, inference_argsNone): # 对 batch 逐条或整体推理返回 Iterable[PredictionResult] ...若管道按 Key 对样本分组且不同 Key 使用不同模型可使用KeyedModelHandler配合KeyModelPathMappingbase.py可一次性将一组 Key 的模型更新到新路径其字段keys、update_path、model_id分别对应受影响 Key 列表、新模型路径与指标中使用的模型标识。实战建议与验证路径优先使用 RunInference 而非手写 ParDo 推理它内置了批量缓冲BatchElements、线程/进程间模型共享、指标收集默认命名空间RunInference与自动模型刷新等能力避免重复造轮子。从示例与测试入手验证接入可参考仓库sdks/python/apache_beam/ml/inference下各框架的*_test.py与*_it_test.py如 pytorch_inference_test.py、sklearn_inference_test.py、tensorflow_inference_test.py它们展示了各ModelHandler的构造参数与预期输出是核对参数含义最直接的依据。流式模型更新部署要点模型文件上传采用唯一命名、保证首份模型先于管道启动就绪、仅在流式管道中使用WatchFilePattern。监控与运维利用RunInference指标model_name参与指标标识追踪当前生效的模型版本与推理耗时如需在多步骤间复用同一模型设置model_identifier以避免重复加载。相关延伸阅读仓库内文档38_ai_runinference.mdRunInference 细节与示例、41_ai_model_refresh.md自动模型刷新、39_ai_runinference_metrics.md推理指标、42_ai_custom_inference.md自定义推理以及 43_ai_llm_inference.mdLLM 推理。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐vscode-graphql-syntax 语法高亮演进全解TextMate 语法、嵌入式注入与 GraphQL 作用域设计vscode graphql syntax 语法高亮演进全解TextMate 语法、嵌入式注入与 GraphQL 作用域设计 GraphiQL 生态中的 vs大数据批处理流处理数据工程GreptimeDB 元数据单事务更新设计表元数据键模型与一致性保障原理GreptimeDB 元数据单事务更新设计表元数据键模型与一致性保障原理 导读 本文基于 GreptimeDB 的 元数据事务 RFC https://lin大数据批处理流处理数据工程一份搞定的 eSearch 全能屏幕工具指南截屏、离线 OCR 与录屏的完整安装教程一份搞定的 eSearch 全能屏幕工具指南截屏、离线 OCR 与录屏的完整安装教程 eSearch 是一款跨平台屏幕工具把截屏、离线 OCR、搜索翻译、以桌面应用OCR屏幕录制视频处理图像处理上一篇DINOv3模型压缩技术从ViT-7B到ConvNeXt Tiny的轻量化策略下一篇Fief源码解析深入理解认证系统的底层实现创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
阅读完成 · 觉得有帮助?