【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载Apache Beam 2.40.0发布于 2022-06-25是本项目仓库中一个具有里程碑意义的版本它首次引入了框架无关的机器学习推理 TransformRunInference官方文档位于 beam-2.40.0.md并为 Go SDK 带来基于泛型的注册函数机制、自检查点 Splittable DoFnSDF等一批重要能力。本文以该版本官方发布说明为骨架结合仓库内 sdks/python/apache_beam/ml/inference 与 sdks/go 的源码实现逐条还原 2.40.0 的亮点、I/O 更新、破坏性变更与已知问题帮助读者在升级或使用该版本时准确把握其能力边界。RunInference2.40.0 的核心亮点2.40.0 发布说明中最重要的新功能是RunInference API——一个与具体机器学习框架解耦的推理 Transform。它定义在 Python SDK 的 sdks/python/apache_beam/ml/inference/base.py 中是一个标准的beam.PTransform接收一个PCollection的样本examples/features输出一个包含原始输入与推理结果的PCollection[PredictionResult]。在该版本中RunInference原生支持PyTorch与Scikit-learn两个框架。为什么需要 RunInference在 RunInference 出现之前在 Beam 流水线中接入机器学习模型需要手工处理模型加载、跨线程模型共享、元素批处理batching与指标采集等重复性工作。而RunInference的模块化设计把这些横切关注点收敛到了ModelHandler抽象中load_model()负责加载并初始化模型run_inference(batch, model, inference_args)对一批样本执行推理get_num_bytes(batch)估算批数据字节数供自动批处理决策使用get_metrics_namespace()为采集到的指标提供命名空间默认值为RunInference。在 base.py 的ModelHandler基类定义中可以看到用户只要针对自己的框架实现load_model与run_inference两个核心方法即可接入任意机器学习库——这正是该 API 号称「framework agnostic」的实现基础。底层执行流程从 RunInference.expand 的实现可以看到该 Transform 的执行流水线包含以下几个阶段预处理Preprocess通过ModelHandler.with_preprocess_fn注册的预处理函数按序应用到每个元素上对应BeamML_RunInference_Preprocess步骤批处理Batch借助beam.BatchElements对元素进行分批批次大小由ModelHandler.batch_elements_kwargs()配置这是提升 GPU/批量推理效率的关键核心推理RunInference_RunInferenceDoFn在 ParDo 中执行真正的模型推理并通过共享机制shared与multi_process_shared让同一模型被多个线程复用避免重复加载后处理Postprocesswith_postprocess_fn注册的后处理函数对推理结果做最终整理。PredictionResult是一个NamedTuple包含example输入样本、inference推理输出与model_id用于区分模型版本的标识三个字段base.py#L82-L97。PyTorch 与 Scikit-learn 的开箱即用2.40.0 为两个框架提供了官方ModelHandler实现PyTorchpytorch_inference.py 提供PytorchModelHandlerTensor与PytorchModelHandlerKeyedTensor。其_load_model实现pytorch_inference.py#L82-L100支持两种模型加载方式(state_dict_path, model_class, model_params)组合或torch_script_model_pathTorchScript 序列化模型同时会校验参数配对例如只传state_dict_path而不传model_class会抛出RuntimeError。代码还包含 GPU 可用性降级逻辑当指定deviceGPU但torch.cuda.is_available()为假时自动回退到 CPU 并打印告警日志。Scikit-learnsklearn_inference.py 提供SklearnModelHandlerNumpy与SklearnModelHandlerPandas通过ModelFileType枚举区分pickle与joblib两种模型序列化格式默认的 Numpy 推理函数会先用numpy.stack将批向量化再调用model.predictsklearn_inference.py#L73-L80。注意joblib是可选依赖未安装时加载 joblib 格式模型会抛出ImportError。官方示例PyTorch 图像分类发布说明中提到的示例位于 sdks/python/apache_beam/examples/inference/pytorch_image_classification.py。该示例展示了 RunInference 的完整用法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_handler KeyedModelHandler( PytorchModelHandlerTensor( state_dict_pathknown_args.model_state_dict_path, model_classmodel_class, # 默认为 mobilenet_v2 model_paramsmodel_params, # {num_classes: 1000} devicedevice, # CPU 或 GPU min_batch_size10, max_batch_size100)).with_preprocess_fn(preprocess).with_postprocess_fn(postprocess) predictions ( filename_value_pair | PyTorchRunInference RunInference(model_handler))关键点由于输入是带文件名的键值对keyed inputs示例用KeyedModelHandler包装了PytorchModelHandlerTensorwith_preprocess_fn(preprocess)负责把图片文件读入并做 Resize(224×224)、ToTensor、归一化with_postprocess_fn(postprocess)用torch.argmax取出预测类别并把(filename, prediction)序列化为文本行最终由WriteToText写出命令行参数包括--input图片文件名列表、--output预测输出路径、--model_state_dict_path模型权重、--images_dir可选图片所在目录。值得注意的是仓库中sdks/python/apache_beam/ml/inference目录还包含tensorflow_inference.py、onnx_inference.py、xgboost_inference.py、vertex_ai_inference.py、tensorrt_inference.py、huggingface_inference.py等更多框架实现与配套测试这反映了 RunInference 在 2.40.0 之后的持续扩展方向。I/O 更新HCatalogIO 升级到 Hive 3.1.3在 I/O 方面2.40.0 将HCatalogIO依赖的 Hive 从旧版本升级到3.1.3对应 Issue-19554。这一变更对使用 HCatalog 读写 Hive 表仓库的用户意味着构建与运行时会解析到更新版本的 Hive 类库带来 bug 修复与兼容性改进用户仍然可以自行提供自己版本的 Hive即在构建/运行时显式覆盖依赖版本说明升级并未强制锁定 Hive 版本留出了向后兼容的弹性。Go SDK 的 2.40.0 演进2.40.0 对 Go SDK 的改动最为密集且全部围绕运行效率与流式处理能力展开。泛型注册函数优化 DoFn 执行2.40.0 起 Go SDK 用户可以使用泛型注册函数generic registration functions来优化 DoFn 执行BEAM-14347。在仓库的 sdks/go/pkg/beam/register 包中可以看到大量基于 Go 泛型的注册辅助函数例如registerStartBundle0x0FuncAndMakeStructWrapper、registerStartBundle1x1FuncAndMakeStructWrapper[I0, R0 any]等register.go它们把形如func(I0) R0的用户函数包装成 runner 可直接调用的reflectx.Func结构体包装器并完成类型注册。泛型带来的收益是注册代码模板化DoFn 的反射调用路径更直接、开销更低。自检查点 SDF 与 textio 全面 SDF 化自检查点 Splittable DoFnBEAM-11104Go SDK 用户现在可以编写自检查点self-checkpointing的 Splittable DoFn用于读取流式数据源。SDF 是 Beam 中支持动态分片dynamic splitting与检查点的关键抽象此项能力让 Go SDK 能更好地支撑有界/无界流式场景。textio Read 全面迁移到 SDFBEAM-14489textio.go 中Read、ReadAll等入口已经改为基于 Splittable DoFn 实现原先的ReadSdf、ReadAllSdf变体被标记为Deprecated并明确注释「Use Read instead, which has been migrated to use this SDF implementation」。Pipeline drain 支持已测试BEAM-11106drain排空优雅停止能力在 Go SDK 上通过测试验证。Worker Status 增强BEAM-13829Go SDK 的 Worker Status 页面现在可以展示堆使用量heap usage、侧输入缓存统计sideinput cache stats与活动 ProcessBundle 统计方便在运行时观察 worker 健康状况。Python 序列化库锁定 dill 0.3.1.12.40.0 将 Python SDK 的序列化pickling库固定为dill0.3.1.1BEAM-11167。dill 是 Beam Python SDK 对闭包、lambda 等复杂 Python 对象进行序列化分发用于跨进程/跨语言执行的基础库版本锁定意味着依赖解析结果可复现同时用户自定义依赖时需要注意与 dill 0.3.1.1 的兼容性。破坏性变更清单升级到 2.40.0 前需要评估以下行为变化变更项影响范围说明Go SDK 最低要求 Go 1.18Go为支持泛型generics编译 Go 流水线的最低 Go 版本提升到 1.18BEAM-14347synthetic.SourceConfig字段类型改为int64Go由int改为int64以更好地与 Flink 在 Schema 中使用的 Logical Types 兼容BEAM-14173默认 coder 变更通用使用BoundedSourceAsSDFWrapperFn与UnboundedSourceAsSDFWrapper时默认 coder 改为对数据源进行压缩compress sources其中int到int64的字段类型变化意味着任何对SourceConfig字段做显式类型断言的用户代码需要同步调整。Bugfixes 与已知问题本次修复Java expansion service 文件暂存修复BEAM-14160允许扩展服务将特定文件进行 staging解决了此前文件无法正确上传的问题。Elasticsearch 连接修复BEAM-14000修复了同时使用 SSL 与用户名/密码认证时 Elasticsearch 连接失败的问题。升级前必须了解的已知问题Pythonbeam.FlatMap与内置函数冲突当用sum、len等内置函数构造beam.FlatMap时会抛出AttributeError: builtin_function_or_method object has no attribute __func__对应 GitHub Issue #22091。绕过方式是在外层包一层用户自定义函数而非直接传入内置函数。JavaBigQueryIO.Write时间戳越界异常当 sink 空闲触发 idle timeout或使用动态目标dynamic destinations时某个表空闲过久写入操作可能尝试输出超出最大时间戳范围的记录报错形如Cannot output with timestamp 294247-01-10T04:00:54.776Z ...此时作业将无法被 drain。官方说明指出该问题已在 2.41 版本修复2.40.0 用户应评估自身写入空闲场景并规划升级窗口。版本下载与完整变更记录本版本的下载入口为官网 download 页面的 2.40.0 条目发布日期 2022-06-25完整变更列表见 GitHub Releases 的 v2.40.0 tag。2.40.0 的贡献者名单由git shortlog统计共包含约 80 位社区开发者。小结Apache Beam 2.40.0 是一次「推理 流式能力」双重驱动的版本Python 侧以 RunInference 开启了 Beam 对机器学习推理的一体化支持PyTorch/Scikit-learn 开箱即用且ModelHandler抽象为后续 TensorFlow、XGBoost、ONNX 等框架的接入铺平了道路Go 侧则借助 Go 1.18 泛型重构了注册机制并将 textio 全面 SDF 化以强化流式读取。与此同时SourceConfig字段类型变更、Go 版本要求提升以及两个已知问题也为计划升级的生产环境划定了明确的检查清单。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam 2.40.0 版本深度解析RunInference API、HCatalogIO 升级与 Go SDK 泛型时代Apache Beam 2.40.0 版本深度解析RunInference API、HCatalogIO 升级与 Go SDK 泛型时代 Apache Bea大数据批处理流处理数据工程Apache Beam 2.39.0 版本解析JMS 动态 Topic 写入、PulsarIO 新连接器与 Go SDK 演进Apache Beam 2.39.0 版本解析JMS 动态 Topic 写入、PulsarIO 新连接器与 Go SDK 演进 Apache Beam 2.3大数据批处理流处理数据工程Apache Beam 2.58.0 版本发布深度解读Solace I/O、RunInference 模型共享与 Storage Write API 改进Apache Beam 2.58.0 版本发布深度解读Solace I/O、RunInference 模型共享与 Storage Write API 改进 本大数据批处理流处理数据工程上一篇Headless Chrome Crawler终极指南10个高效爬取技巧让数据采集更简单下一篇graphify 并行语义抽取的 Step B2 全量子代理分发协议单条消息并行调度、general-purpose 落盘契约与 chunk 路径规范创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考