
批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载Apache Beam 2.51.02023 年 10 月 3 日发布是一次聚焦机器学习推理能力与跨语言 SDK 健壮性的重要版本。本文基于 官方发布公告结合本仓库sdks/python、sdks/java、sdks/go下的真实源码与测试逐条解析该版本的新特性、破坏性变更、Bug 修复、安全修复与已知问题帮助正在升级或维护 Beam 管道的开发者评估影响范围并顺利完成迁移。一、版本概览与发布背景2.51.0 是 Apache Beam 在 2023 年 10 月发布的稳定版本其发布公告位于 beam-2.51.0.md主要围绕三条主线展开Python 机器学习推理RunInference能力增强支持在同一个 Transform 内通过KeyedModelHandler加载多个模型并允许VertexAIModelHandlerJSON透传inference_args到 Vertex AI 端点。平台级工程质量提升新增对用户管道运行mypy静态检查的支持同时修复 GCS 连接器异常链、流式插入异常处理、跨语言 Bigtable sink 时间戳等多个历史问题。跨 SDK 的破坏性变更涉及 Java Beam SQL移除 fastjson、表属性改为 Jackson ObjectNode、Python 容器镜像移除 TensorFlow、Go SDKparquetio.Write签名简化、以及BeamSqlSeekableTable.setUp接口调整。注意当前仓库为 Beam 主分支快照本文涉及的具体文件路径与实现细节以仓库实际内容为准2.51.0 的详细逐条变更记录可参见发布里程碑的 release notes。二、新特性与改进详解1. Python RunInference 支持 KeyedModelHandler 多模型加载公告要点Python 的RunInference现在支持在同一个 Transform 中通过KeyedModelHandler加载多个模型对应 issue #27628。KeyedModelHandler的完整实现位于 sdks/python/apache_beam/ml/inference/base.py。从源码可以看出它是一个泛型类class KeyedModelHandler(Generic[KeyT, ExampleT, PredictionT, ModelT], ModelHandler[Tuple[KeyT, ExampleT], Tuple[KeyT, PredictionT], Union[ModelT, _ModelManager]]): def __init__( self, unkeyed: Union[ModelHandler[ExampleT, PredictionT, ModelT], List[KeyModelMapping[KeyT, ExampleT, PredictionT, ModelT]]], max_models_per_worker_hint: Optional[int] None):其核心设计思路是输入输出均带 Key如果原始模型配合RunInference将PCollection[E]转为PCollection[P]那么KeyedModelHandler则处理PCollection[Tuple[K, E]]并产出PCollection[Tuple[K, P]]通过 Key 将输入与输出关联起来。两种构造方式传入单个无 Key 的ModelHandler_single_model not isinstance(unkeyed, list)此时它只是把 Key 与推理结果重新 zip 起来传入KeyModelMapping列表将不同的 Key 区间映射到不同的ModelHandler。KeyModelMapping是一个 dataclassbase.py用法示例k1 [k1, k2, k3] k2 [k4, k5] KeyedModelHandler([KeyModelMapping(k1, mh1), KeyModelMapping(k2, mh2)])底层实现多模型场景run_inference先按 Key 将 batch 分组batch_by_key、key_by_id再为每个模型 cohort 通过multi_process_shared.MultiProcessShared获取共享模型实例逐个调用底层mh.run_inference最后按 Key 顺序组装预测结果base.py。使用要点与内存管理多个模型可能同时驻留内存加载过多大模型可能触发 OOM为此KeyedModelHandler提供max_models_per_worker_hint参数提示 runner 每个 worker 进程同一时刻最多加载多少个模型base.py。多模型场景下各底层 ModelHandler 的 batching kwargs、resource hints、_env_vars会被忽略源码中会打印 warning如需自定义需覆写KeyedModelHandler.batch_elements_kwargs()等方法。KeyedModelHandler 还支持 Automatic Model Refresh通过update_model_paths方法配合KeyModelPathMapping侧输入可在不重启流式管道的情况下更新模型路径base.py。指标采集多模型时既会按 cohort 聚合无前缀也会按 cohort 输出cohort_key-metric_name形式的指标模型更新后变为cohort_key-model id-metric_name。在base_test.py中可以看到针对KeyedModelHandler的覆盖性测试如test_run_inference_impl_inference_args、test_unexpected_inference_args_passed等验证了inference_args的透传与校验行为sdks/python/apache_beam/ml/inference/base_test.py。2. VertexAIModelHandlerJSON 支持透传 inference_args公告要点Python 的VertexAIModelHandlerJSON现在支持传入inference_args这些参数会作为parameters透传给 Vertex AI 端点。实现位于 sdks/python/apache_beam/ml/inference/vertex_ai_inference.pyretry.with_exponential_backoff( num_retries5, retry_filter_retry_on_appropriate_gcp_error) def get_request( self, batch: Sequence[Any], model: aiplatform.Endpoint, throttle_delay_secs: int, inference_args: Optional[Dict[str, Any]]): ... prediction model.predict( instanceslist(batch), parametersinference_args)关键点远程推理模式与其它ModelHandler在 worker 上加载模型不同VertexAIModelHandlerJSON是向 Vertex AI 端点发起远程查询行为更接近管道中段的 IO。inference_args 的流向run_inference→get_request→aiplatform.Endpoint.predict(instances..., parametersinference_args)即字典会被整体作为预测请求的parameters字段发送到端点供模型服务端按需解释如 top_k、temperature 等。请求限制公共 Vertex AI 端点单次请求上限 1.5 MB如需更大请求需配置 Compute Engine network 并设置privateTrue此时必须同时提供network否则构造函数抛出ValueError见 vertex_ai_inference.py。客户端限流与重试源码内置了AdaptiveThrottler窗口 1ms、bucket 1ms、overload_ratio2实现客户端自适应限流并通过_retry_on_appropriate_gcp_error仅对 5xx 与 429TooManyRequests进行指数退避重试num_retries5。相关重试过滤器行为在 vertex_ai_inference_test.py 中有单元测试覆盖。构造函数参数endpoint_id端点数值 ID、projectGCP 项目、location区域、experiment实验标签、network私有端点 VPC 全名、private是否为私有端点、以及 batching 相关的min_batch_size、max_batch_size、max_batch_duration_secs。3. 支持对用户管道运行 mypy 静态检查公告要点新增支持对用户管道运行mypy对应 issue #27906。在 sdks/python/setup.py 中新增了mypy这一 distutilsCommand子类class mypy(Command): ... def run(self): args [mypy, self.get_project_path()] result subprocess.call(args) if result ! 0: raise DistutilsError(mypy exited with status %d % result)它先运行egg_info与就地构建扩展再对构建出的项目路径执行mypy并在返回码非 0 时以DistutilsError终止。这样 Beam Python SDK 自身的类型检查能力可以平滑复用到用户的管道代码上。仓库同时在 mypy.ini、pyproject.toml 与 tox.ini 中维护了 mypy 配置使用者可以参考其中的 strict 开关、忽略规则与插件设置来对齐官方检查标准。三、破坏性变更逐项评估升级必读1. Beam SQL 移除 fastjson表属性改为 Jackson ObjectNodeJava公告要点Beam SQL 移除了 fastjson 库依赖表Table属性改为基于 JacksonObjectNode对应 issue #24154。该变更的影响面覆盖整个sdks/java/extensions/sql模块。从当前仓库源码看Beam SQL 的表属性定义已全面基于 Jackson例如 TableUtils.java、Table.java 以及各TableProvider如 KafkaTableProvider.java、BigQueryTableProvider.java、TextTableProvider.java 等。迁移建议若你的自定义TableProvider或管道代码直接解析/构造表属性 JSON需从 fastjson 的对象模型切换到 JacksonObjectNode关注Table对象properties()的返回类型变化所有基于旧类型的序列化/反序列化代码都应同步调整该变更同时移除了 fastjson 这一第三方依赖有助于缩小依赖面与潜在安全暴露。2. Python 容器镜像移除 TensorFlow公告要点Beam Python 容器镜像中移除了 TensorFlow对应 PR #28424。受影响的用户包括依赖官方 Python 容器内置 TensorFlow 运行 RunInference 的管道。若你受到此变更影响可在 issue #20605 中反馈。迁移建议改用自定义容器镜像在镜像中显式安装所需版本的tensorflow并通过 Dataflow / Runner v2 的容器配置--experimentsuse_runner_v2及自定义 worker 镜像参数指定该镜像。3. Go SDKparquetio.Write 移除 reflect.Type 参数公告要点parquetio.Write移除了t reflect.Type参数元素类型改为从输入 PCollection 推导对应 issue #28490。当前仓库 sdks/go/pkg/beam/io/parquetio/parquetio.go 中的签名已变为// Write writes a PCollectionparquetStruct to .parquet file. // Write expects elements of a struct type with parquet tags func Write(s beam.Scope, filename string, col beam.PCollection) { t : col.Type().Type() s s.Scope(parquetio.Write) filesystem.ValidateScheme(filename) pre : beam.AddFixedKey(s, col) post : beam.GroupByKey(s, pre) beam.ParDo0(s, parquetWriteFn{Filename: filename, Type: beam.EncodedType{T: t}}, post) }旧代码形如parquetio.Write(s, filename, t, col)的调用需要删掉中间的t改为parquetio.Write(s, filename, col)。元素类型仍要求为带parquettag 的 struct如文档注释中Student示例类型信息通过col.Type().Type()在内部自动获取写入时由parquetWriteFn.ProcessElement借助writer.NewParquetWriterFromWriter完成。4. BeamSqlSeekableTable.setUp 新增 joinSubsetType 参数公告要点BeamSqlSeekableTable.setUp重构新增参数joinSubsetType对应 issue #28283。当前 BeamSqlSeekableTable.java 中的默认方法签名为/** * param joinSubsetType joining subset schema */ default void setUp(Schema joinSubsetType) {}该参数携带 join 子集的 schema 信息便于可查找表在 lookup join 场景下按需准备数据。任何实现了BeamSqlSeekableTable接口并覆写setUp的自定义表类都需要同步更新方法签名。相关调用方如 BeamSideInputLookupJoinRel.java、BeamJoinTransforms.java均已按新签名调整测试见 BeamSideInputLookupJoinRelTest.java。四、Bug 修复解析1. GCS 连接器异常链修复Python修复了 GCS 连接器中的异常链exception chaining问题对应 issue #26769基于google-cloud-storage的 BlobReader/BlobWriter 与重试机制中异常抛出时对原始异常的保留与透传使上层管道能够定位到真正的根因而不是被包装后的通用错误掩盖。2. 流式插入异常处理修复Python公告要点修复了流式插入的异常处理GoogleAPICallErrors现在会按照重试策略进行重试并在适当情况下路由到失败行failed rows而非直接导致管道失败对应 issue #21080。这意味着使用 BigQuery 流式插入WriteToBigQuery的WriteDisposition/插入流模式的管道面对瞬时GoogleAPICallErrors时具备更强的容错性可重试错误进入重试路径不可恢复错误被定向到 DLQ / failed rows 集合避免整个管道因单批数据失败而崩溃。3. 跨语言 Bigtable sink 时间戳处理修复Python修复了 Python SDK 的跨语言 Bigtable sink 中未显式设置时间戳的记录被错误处理的问题对应 issue #28632。此前缺少显式时间戳的记录在通过 Cross-Language Transform 写入 Bigtable 时可能产生非预期行为修复后这类记录按缺省语义正确处理。五、安全修复2.51.0 对运行环境依赖做了同步升级以消除已知 CVE类别修复内容涉及 CVEPython 容器更新容器内依赖CVE-2021-30473、CVE-2021-30474、CVE-2021-30475commons-compress 相关、CVE-2020-36130、CVE-2020-36131、CVE-2020-36133、CVE-2020-36135Go 工具链使用 go 1.21.1 构建CVE-2023-39320Go 标准库如果你的管道使用官方发布的基础容器镜像升级到 2.51.0 并重新拉取对应 tag 即可获得上述安全修复自建镜像的用户应同步更新容器内的 Python/Go 依赖版本。六、已知问题与规避方案已知问题使用 BigQuery Storage Read API 的 Python 管道必须将fastavro依赖固定到1.8.3 或更早版本对应 issue #28811。依据仓库中 BigQuery Storage Read API 的 Avro 反序列化路径直接依赖fastavro——如 bigquery.py 的_AvroRecordIteratorfastavro.parse_schemafastavro.schemaless_reader、bigquery_tools.py 的fastavro.write.Writer使用。规避方法是在管道依赖声明中显式固定fastavro1.8.3或在requirements.txt/ 自定义容器镜像中pip install fastavro1.8.3直至官方在后续版本中解除该约束。七、升级行动清单综合上述变更从 2.50.x 升级到 2.51.0 建议按以下顺序检查Java Beam SQL排查自定义 TableProvider / 表属性解析代码中的 fastjson 依赖迁移到 JacksonObjectNode更新所有BeamSqlSeekableTable自定义实现的setUp(Schema joinSubsetType)签名。Python若依赖官方容器内嵌 TensorFlow改用自定义镜像并显式安装 TensorFlow使用 Storage Read API 时固定fastavro1.8.3尝试用python setup.py mypy或仓库中的 mypy.ini 配置对你的管道做静态检查。Go将parquetio.Write(s, filename, t, col)改为parquetio.Write(s, filename, col)元素类型由 PCollection 自动推导。依赖与安全确认容器镜像已更新至 2.51.0含 Python 依赖 CVE 修复与 Go 1.21.1 构建产物。验证回归重点回归 RunInference尤其多模型 KeyedModelHandler、BigQuery 流式插入失败行路由、跨语言 Bigtable sink 三条链路。通过上述清单你可以将 2.51.0 的核心变化快速映射到自己的管道代码并在升级过程中精准定位需要修改的位置。详细的分项变更记录含 issue 链接与具体 commit可对照官方发布里程碑的 release notes 逐条核对。赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐Apache Beam 2.51.0 版本全解析多模型 RunInference、Vertex AI 推理增强与破坏性变更指南Apache Beam 2.51.0 版本全解析多模型 RunInference、Vertex AI 推理增强与破坏性变更指南 Apache Beam 2.5Apache Beam 2.51.0 版本发布解读Python RunInference 多模型推理、Vertex AI 参数透传与破坏性变更一览Apache Beam 2.51.0 版本发布解读Python RunInference 多模型推理、Vertex AI 参数透传与破坏性变更一览 Apach大数据批处理流处理数据工程Apache Beam 2.46.0 版本深度解析RunInference 模型服务增强、I/O 能力扩展与破坏性变更指南Apache Beam 2.46.0 版本深度解析RunInference 模型服务增强、I/O 能力扩展与破坏性变更指南 Apache Beam 2.46.大数据批处理流处理数据工程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考