- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
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 侧约束收紧;
- 弃用:MongoDb
withKeepAlive(因 Mongo 驱动中已弃用)。
依赖升级与变更
Java 依赖升级
2.11.0 对 Java 生态依赖进行了系统性升级,直接影响到使用这些组件的 I/O 连接器(如 gRPC 相关连接器、Bigtable、Spanner、Pub/Sub、GCS 等)。核心升级包括:
| 类别 | 组件 | 新版本 |
|---|---|---|
| 解析 | antlr / antlr_runtime | 4.7 |
| gRPC 全家桶 | grpc_all / grpc_auth / grpc_core / grpc_pubsub_v1 / grpc_protobuf / grpc_protobuf_lite / grpc_netty / grpc_stub | 1.17.1 |
| Netty | netty_handler / netty_transport_native_epoll | 4.1.30.Final |
| netty_tcnative_boringssl_static | 2.0.17.Final | |
| Google Cloud | google_api_common / google_auth_library_credentials / google_auth_library_oauth2_http | 1.7.0 / 0.12.0 / 0.12.0 |
| google_cloud_core / google_cloud_core_grpc | 1.61.0 | |
| google_api_services_dataflow | v1b3-rev20190126-1.27.0 | |
| google_cloud_bigquery_storage / google_cloud_bigquery_storage_proto | 0.79.0-alpha / 0.44.0 | |
| google_cloud_spanner / proto_google_cloud_spanner_admin_database_v1 | 1.6.0 / 1.6.0 | |
| bigtable_client_core | 1.8.0 | |
| gax_grpc | 1.38.0 | |
| 数据访问 | bigdataoss_gcsio / bigdataoss_util | 1.9.16 |
| cassandra-driver-core / cassandra-driver-mapping | 3.6.0 | |
| 压缩 | commons-compress | 1.18 |
| zstd_jni | 1.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.27google-cloud-core==0.28.1google-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)实际运行两个消费者:
- 主消费者(main consumer):真正从 Kafka 读取数据;
- 次级偏移量消费者(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 | 默认关闭 |
withMinRowCountForPageSizeCheck | page 大小检查前最少缓冲行数,大行场景可调低(如 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 有两项针对性优化:
- GroupByKey 非合并窗口优化:当窗口无需合并(non-merging windows,如 FixedWindows/SlidingWindows 之外的固定窗口场景)时,GroupByKey 的调度与数据组织得到优化,减少了不必要的窗口合并开销;
- 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 的升级评估,可遵循以下要点:
- 核对依赖:Java 侧 gRPC 1.17.1、Netty 4.1.30、Google Cloud 组件版本均有提升,确认与自身项目无版本冲突;
- 利用新 I/O 能力:安全 Kafka 集群务必配置
withOffsetConsumerConfigOverrides;Parquet 输出按存储与查询成本选择压缩编解码器;BigQuery/GCS 场景通过 KMS 密钥满足加密合规; - Python 3 迁移:可开始用 DirectRunner/DataflowRunner 验证 Python 3 实验性支持,但生产环境需评估实验状态的风险;
- Runner 调优:Spark 场景可使用
bundleSize控制分裂粒度,Flink 流式作业可利用 savepoint 升级能力规划滚动升级方案; - 清理弃用 API:移除 MongoDB
withKeepAlive调用。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam 2.18.0 版本全解析:Spark Structured Streaming Runner、SQS/RabbitMQ I/O 与 SQL 能力升级
Apache Beam 2.18.0 版本全解析:Spark Structured Streaming Runner、SQS/RabbitMQ I/O 与 SQ
大数据批处理流处理数据工程WSABuilds 完整指南:3 步在 Windows 上跑起安卓应用
WSABuilds 完整指南:3 步在 Windows 上跑起安卓应用 想玩的某款安卓游戏,Windows 上翻来覆去只有个只带亚马逊应用商店的 WSA,Goo
大数据批处理流处理数据工程Apache Beam 2.15.0 发布解读:I/O 增强、SQL ParquetTable 与 Runner 改进全解析
Apache Beam 2.15.0 发布解读:I/O 增强、SQL ParquetTable 与 Runner 改进全解析 Apache Beam 2.15.
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考