☰
Apache Beam 2.11.0 版本解析:依赖升级、新 I/O 能力与运行时改进全指南
2026/10/9 3:20:33 网站建设 项目流程
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载

Apache Beam 2.11.0 是该项目于 2019 年发布的一个重要里程碑版本,聚焦于"改进与新增功能"双线推进。本文以官方发布博客 beam-2.11.0.md 为核心骨架,结合当前仓库源码逐项拆解该版本的依赖升级清单、I/O 能力增强(Kafka 偏移量消费者、Parquet 压缩编解码器、BigQuery KMS 加密、GCS KMS 拷贝)、新特性(Python 3 实验性支持、ZStandard 压缩、CombineFn.compact、Spark/Flink Runner 优化)及弃用项,帮助读者快速评估升级收益与迁移影响。

版本概览

Apache Beam 2.11.0 于 2019-02-26 发布,官方公告见仓库中的 beam-2.11.0.md。该版本同时包含改进与新增功能,主要亮点集中在:

  • I/O 能力增强:跨语言变换的 Portable Flink Runner 支持、GCS 拷贝的 Cloud KMS 支持、KafkaIO 偏移量消费者参数、ParquetIO 写入压缩编解码器、BigQuery KMS 密钥传递;
  • 新特性:Python 3(实验性)DirectRunner/DataflowRunner 支持、Java SDK 的 ZStandard 压缩、Python CombineFn.compact、Spark Runner GroupByKey 非合并窗口优化与 bundleSize 参数、Flink Runner 可移植 Runner savepoint/升级支持;
  • 依赖大规模升级:Java 侧 grpc、netty、google 生态组件批量升级,Python 侧约束收紧;
  • 弃用:MongoDbwithKeepAlive(因 Mongo 驱动中已弃用)。

依赖升级与变更

Java 依赖升级

2.11.0 对 Java 生态依赖进行了系统性升级,直接影响到使用这些组件的 I/O 连接器(如 gRPC 相关连接器、Bigtable、Spanner、Pub/Sub、GCS 等)。核心升级包括:

类别组件新版本
解析antlr / antlr_runtime4.7
gRPC 全家桶grpc_all / grpc_auth / grpc_core / grpc_pubsub_v1 / grpc_protobuf / grpc_protobuf_lite / grpc_netty / grpc_stub1.17.1
Nettynetty_handler / netty_transport_native_epoll4.1.30.Final
netty_tcnative_boringssl_static2.0.17.Final
Google Cloudgoogle_api_common / google_auth_library_credentials / google_auth_library_oauth2_http1.7.0 / 0.12.0 / 0.12.0
google_cloud_core / google_cloud_core_grpc1.61.0
google_api_services_dataflowv1b3-rev20190126-1.27.0
google_cloud_bigquery_storage / google_cloud_bigquery_storage_proto0.79.0-alpha / 0.44.0
google_cloud_spanner / proto_google_cloud_spanner_admin_database_v11.6.0 / 1.6.0
bigtable_client_core1.8.0
gax_grpc1.38.0
数据访问bigdataoss_gcsio / bigdataoss_util1.9.16
cassandra-driver-core / cassandra-driver-mapping3.6.0
压缩commons-compress1.18
zstd_jni1.3.8-3

提示:表格中各组件的新版本号与官方发布说明一致。对于使用 Google Cloud 连接器的任务,升级到 2.11.0 时建议核对自身依赖树,避免与上表版本产生冲突。

Python 依赖变更

Python SDK 在 2.11.0 中收紧了部分依赖版本约束,其中带有python_version < "3.0"条件的项仅作用于 Python 2 环境:

  • futures>=3.2.0,<4.0.0; python_version < "3.0"(Python 2 下并发原语依赖)
  • pyvcf>=0.6.8,<0.7.0; python_version < "3.0"
  • google-apitools>=0.5.26,<0.5.27
  • google-cloud-core==0.28.1
  • google-cloud-bigtable==0.31.1

这与该版本同步引入的 Python 3 实验性支持相呼应:Python 2 环境继续通过条件依赖保持兼容,而 Python 3 环境则开始获得独立运行能力。

I/O 能力增强

Portable Flink Runner 跨语言变换支持

2.11.0 中 Portable Flink Runner 开始支持运行跨语言(cross-language)变换,即在同一 Pipeline 中混合使用不同 SDK(如 Java 与 Python)编写的变换。该能力建立在 Beam 的可移植性架构之上,使 Flink 作为执行后端时可以调度来自多语言 SDK 的扩展服务。这意味着用户可以利用 Python/Java 各自的生态(例如使用 Python 编写的数据处理逻辑配合 Java 的成熟 I/O 连接器),由 Flink Runner 统一调度执行。

GCS 拷贝的 Cloud KMS 支持

该版本为 GCS 拷贝操作增加了 Cloud KMS 支持,允许在跨存储桶复制对象时使用客户管理的加密密钥(CMEK)。这为数据在 GCS 间的迁移场景提供了加密密钥的可控性,使企业可以在合规要求下统一管理对象加密密钥。

KafkaIO 偏移量消费者配置

2.11.0 为KafkaIO.read()新增了偏移量消费者(offset consumer)相关参数。从当前仓库的 KafkaIO.java 源码看,Kafka 读取后端(ReadFromKafkaDoFn)实际运行两个消费者:

  1. 主消费者(main consumer):真正从 Kafka 读取数据;
  2. 次级偏移量消费者(secondary offset consumer):通过拉取每个分区的latest offset来估算 backlog(积压量),用于推进 watermark 与流量控制。

默认情况下,偏移量消费者继承主消费者的配置,并使用自动生成的group.id。但在安全加固的 Kafka 集群(如启用 SASL/SSL)中,这一默认行为可能失败,运行时会出现如下 WARN 日志:

exception while fetching latest offset for partition {}. will be retried

此时可通过新增的配置 API 为偏移量消费者单独注入配置。仓库中相关方法为 withOffsetConsumerConfigOverrides,用法示例:

KafkaIO.<String, String>read() .withBootstrapServers("broker:9092") .withTopic("my-topic") .withOffsetConsumerConfigOverrides(ImmutableMap.of( "sasl.mechanism", "PLAIN", "security.protocol", "SASL_PLAINTEXT")) ...

同时,仓库还提供了对称的withConsumerConfigUpdates(在默认消费者属性基础上合并更新)与withConsumerConfigOverrides(整体替换主消费者配置)等方法(见 KafkaIO.java#L2841-L2865 与 KafkaIO.java#L3016-L3019),两者配合即可分别治理主消费者与偏移量消费者的配置。

ParquetIO 写入压缩编解码器

2.11.0 允许在ParquetIO.write()中显式设置压缩编解码器。仓库中 ParquetIO.java 的Sink.withCompressionCodec(CompressionCodecName)即对应此能力,默认值为CompressionCodecName.SNAPPY(见 ParquetIO.java#L1069-L1072):

pipeline.apply(...) .apply(ParquetIO.write(FileSystems.matchNewResource("/tmp/out.parquet", false)) .withSchema(schema) .withCompressionCodec(CompressionCodecName.GZIP));

该设置会经由open()中的AvroParquetWriter.Builder传递到底层 Parquet 写入器(ParquetIO.java#L1191-L1200)。写入 Sink 上还提供了系列配套参数,便于在开启压缩的同时精细控制文件布局与内存占用:

方法作用默认值/约束
withCompressionCodec设置压缩编解码器SNAPPY
withConfiguration指定 Hadoop Configuration(Map 或 Configuration)—
withRowGroupSize设置 row-group 大小必须为正整数,否则使用底层默认值
withPageSize设置 page 大小1 MB
withDictionaryEncoding开关字典编码默认开启
withBloomFilterEnabled开关 bloom filter默认关闭
withMinRowCountForPageSizeCheckpage 大小检查前最少缓冲行数,大行场景可调低(如 1)100
withMaxRowCountForPageSizeCheck强制 page 大小检查的最大缓冲行数,防止行大小差异大时缓冲区溢出默认由 Parquet 估算,上限 10000

BigQuery 变换的 KMS 密钥传递

该版本为 BigQuery 变换增加了kms_key支持并传递给 Dataflow 执行。仓库中 BigQueryIO.java 的读(TypedRead、DynamicRead)与写(Write)侧均提供withKmsKey(String)方法(见 BigQueryIO.java#L3568-L3569):

BigQueryIO.writeTableRows() .to("project:dataset.table") .withKmsKey("projects/my-project/locations/global/keyRings/my-ring/cryptoKeys/my-key")

在批量加载实现 BatchLoads.java 中,kmsKey作为加载作业参数被传递,从而在 Dataflow 执行 BigQuery 导入/导出时使用指定的客户管理密钥加密目标表,满足数据静态加密的合规诉求。

新特性与改进

Python 3 实验性支持

2.11.0 为 DirectRunner 与 DataflowRunner 引入了 Python 3 的实验性支持。这是 Beam Python SDK 向 Python 3 迁移进程中的重要一步,意味着上述两个 Runner 已可运行基于 Python 3 编写的 Pipeline,但该能力在当时仍处于实验阶段,官方尚未将其标注为生产级稳定支持。与之配套的 Python 依赖约束(如futures仅在 Python 2 下引入)也体现了双版本共存的过渡策略。

Java SDK 的 ZStandard 压缩支持

该版本为 Java SDK 增加了 ZStandard(zstd)压缩支持。仓库中 ZstdCoder.java 提供了若干工厂方法,可用于组合出带 zstd 压缩的 Coder:

// 包裹内层 coder,使用默认压缩级别 ZstdCoder<T> coder = ZstdCoder.of(innerCoder); // 指定压缩级别 ZstdCoder<T> coder = ZstdCoder.of(innerCoder, level); // 同时指定压缩字典与级别 ZstdCoder<T> coder = ZstdCoder.of(innerCoder, dict, level);

zstd 以高压缩比与较快的压缩/解压速度著称,适合对 PCollection 中间数据或落盘数据进行压缩以节省存储与网络带宽。其依赖zstd_jni也正是上文依赖升级表中被提升至 1.3.8-3 的组件。

Python CombineFn.compact

Python SDK 新增CombineFn.compact,与 Java SDK 的CombineFn.compact行为对齐。仓库中 core.py 的CombineFn基类及其派生实现均声明了compact(accumulator, *args, **kwargs)(如 core.py#L1165),其作用是在合并累加器之前对累加器进行压缩/规约,以减少内存占用和后续合并的开销。典型场景是累加器体积较大、且存在冗余信息可在合并前去重的组合逻辑(如超大规模去重、直方图/摘要统计等),实现该方法可显著降低分布式执行中的传输与存储成本。

Spark Runner 优化

2.11.0 中 Spark Runner 有两项针对性优化:

  1. GroupByKey 非合并窗口优化:当窗口无需合并(non-merging windows,如 FixedWindows/SlidingWindows 之外的固定窗口场景)时,GroupByKey 的调度与数据组织得到优化,减少了不必要的窗口合并开销;
  2. bundleSize 参数:新增bundleSize参数用于控制 Spark source 的分裂(splitting)粒度。仓库中 SparkPipelineOptions.java 定义了该选项,其语义为"若设置了 bundleSize,将用它来分裂 BoundedSource,否则使用默认值";在 SourceRDD.java 中,bundleSize > 0时按该值决定每个分片的目标字节数,否则回退到DEFAULT_BUNDLE_SIZE。合理调大 bundleSize 可减少任务切分数、降低调度开销,调小则可提升并行度,需要结合数据规模与集群资源权衡。

Flink Runner:可移植 Runner savepoint / 升级支持

Flink Runner 的便携式(portable)执行路径新增了 savepoint 与作业升级支持,允许用户在升级作业拓扑或 Beam 版本时,基于 Flink 的 savepoint 机制保留状态、实现无状态丢失的平滑升级。该能力对生产环境中的流式作业尤为重要。

Bugfixes 与弃用

Bugfixes

官方发布说明以"Various bug fixes and performance improvements"概括了该版本的缺陷修复与性能改进清单。具体条目可参考 Apache JIRA 上该版本的 Release Notes(详见原文档中指向的 JIRA 链接),升级时建议结合自身使用到的连接器与 Runner 核对其修复项。

弃用:MongoDb withKeepAlive

2.11.0 弃用了 MongoDB 连接器的withKeepAlive配置,原因是该参数对应的功能已在底层 Mongo 驱动中废弃。使用旧版 API 的用户应迁移到驱动推荐的心跳/保活机制,避免依赖已弃用行为。

贡献者

根据git shortlog统计,2.11.0 版本共有 60 余位贡献者参与,完整名单见官方发布博客 beam-2.11.0.md 的 Contributors 小节,其中包括 Ahmet Altay、Kenneth Knowles、Robert Bradshaw、Tyler Akidau、Reuven Lax 等 Apache Beam PMC 成员与活跃社区开发者。

升级建议小结

针对 2.11.0 的升级评估,可遵循以下要点:

  1. 核对依赖:Java 侧 gRPC 1.17.1、Netty 4.1.30、Google Cloud 组件版本均有提升,确认与自身项目无版本冲突;
  2. 利用新 I/O 能力:安全 Kafka 集群务必配置withOffsetConsumerConfigOverrides;Parquet 输出按存储与查询成本选择压缩编解码器;BigQuery/GCS 场景通过 KMS 密钥满足加密合规;
  3. Python 3 迁移:可开始用 DirectRunner/DataflowRunner 验证 Python 3 实验性支持,但生产环境需评估实验状态的风险;
  4. Runner 调优:Spark 场景可使用bundleSize控制分裂粒度,Flink 流式作业可利用 savepoint 升级能力规划滚动升级方案;
  5. 清理弃用 API:移除 MongoDBwithKeepAlive调用。
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载
上一篇:如何快速诊断nvm-windows问题:自动化日志分析工具终极指南
下一篇:ObjectivePGP高级技巧:如何优化加密性能与安全性

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询