- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
Apache Beam 2.39.0 于 2022 年 5 月 25 日发布,是 Beam 在统一批流编程模型上继续演进的一个重要版本。本篇文章以该版本的官方发布说明为主体,结合当前仓库中的源码实现,深入剖析 JMS 动态 Topic 写入、Apache PulsarIO 连接器、Python DataFrame API 扩展、Go SDK Pipeline Drain、Dataflow 模拟凭据等核心变更的来龙去脉,帮助读者理解这些特性背后的 API 设计与实际用法。读完本文,你将掌握 2.39.0 中新增 I/O 能力的配置方式、破坏性变更的迁移要点,以及如何在自己的管道中落地这些新功能。
本文内容均基于仓库内的发布文档 beam-2.39.0.md 及对应源码展开,仓库当前状态可能与 2.39.0 发布时的代码存在演进差异,源码引用仅用于佐证 API 设计。
版本概览
2.39.0 是 Apache Beam 于 2022 年 5 月 25 日发布的功能版本(官方下载页对应#2390-2022-05-25),主要涵盖四类变更:
- I/Os:JMS 连接器能力大幅增强、新增 Apache PulsarIO、BigQueryIO 流式写入线程优化;
- New Features / Improvements:Flink Scala 2.12 支持、Go SDK Pipeline drain、Python DataFrame API 扩展、Dataflow 模拟凭据、Elasticsearch 8.x 支持、Kinesis 分片感知聚合等;
- Breaking Changes:JmsIO 强制要求
valueMapper、Python Coder 继承约束、FileSystem 新增metadata()抽象方法、Go pipelinex 无用函数移除; - Deprecations:Flink 1.11 与 Python 3.6 停止支持。
I/Os:连接器层的关键升级
JmsIO:任意输入映射为 JMS Message 与动态 Topic 写入
2.39.0 之前,JmsIO 的写入端(JmsIO.Write)只能把元素发送到静态配置的 queue 或 topic。2.39.0 通过 BEAM-16308 引入两项能力:
- 任意类型映射:可以把任何类型的输入元素映射为
javax.jms.Message的任意子类(如TextMessage、ObjectMessage、BytesMessage),不再局限于字符串文本; - 动态 Topic 写入:支持根据输入元素内容在运行时动态决定目标 Topic 名称,实现按数据路由的发布。
两个核心 Mapper 的语义
topicNameMapper:类型为SerializableFunction<EventT, String>,接收一个输入事件对象,返回该事件应该发往的 Topic 名称。它必须与withQueue()、withTopic()三者互斥使用。valueMapper:类型为SerializableBiFunction<EventT, Session, Message>,接收输入事件与一个 JMSSession,返回一个 JMSMessage实例。2.39.0 起该项成为必填配置(破坏性变更,详见下文)。
从 JmsIO.java 的源码看,写入端expand()会做三重校验(对应第 1300-1313 行):
checkArgument( getConnectionFactoryProviderFn() != null, "Either withConnectionFactory() or withConnectionFactoryProviderFn() is required"); checkArgument( getTopicNameMapper() != null || getQueue() != null || getTopic() != null, "Either withTopicNameMapper(topicNameMapper), withQueue(queue), or withTopic(topic) is required"); checkArgument(getValueMapper() != null, "withValueMapper() is required");即:连接工厂必须提供;queue/topic/topicNameMapper 三者必须且只能指定一个;valueMapper必填。底层发送逻辑(第 1407-1413 行)会在每个元素上先调用valueMapper生成 JMS Message,若配置了topicNameMapper则通过session.createTopic(...)创建动态目标:
Message message = spec.getValueMapper().apply(input, session); if (spec.getTopicNameMapper() != null) { destinationToSendTo = session.createTopic(spec.getTopicNameMapper().apply(input)); }动态 Topic 写入示例
源码 Javadoc 给出按公司/员工 ID 路由 Topic 的典型用法:
SerializableFunction<CompanyEvent, String> topicNameMapper = (event -> String.format( "company/%s/employee/%s", event.getCompanyName(), event.getEmployeeId())); pipeline .apply(...) // PCollection<CompanyEvent> .apply(JmsIO.write() .withConnectionFactory(jmsConnectionFactory) .withTopicNameMapper(topicNameMapper) .withValueMapper(valueMapper));其中valueMapper可以把事件序列化为 JMSTextMessage:
SerializableBiFunction<SomeEventObject, Session, Message> valueMapper = (e, s) -> { try { TextMessage msg = s.createTextMessage(); msg.setText(Mapper.MAPPER.toJson(e)); return msg; } catch (JMSException ex) { throw new JmsIOException("Error!!", ex); } };读取端的 MessageMapper
与写入端对应,读取端也支持把任意 JMS Message 映射为自定义 POJO。JmsIO.readMessage()通过withMessageMapper(MessageMapper<T>)完成转换(源码 JmsIO.java 第 106-120 行的 Javadoc 示例):
pipeline.apply(JmsIO.<T>readMessage() .withConnectionFactory(myConnectionFactory) .withQueue("my-queue") .withMessageMapper((MessageMapper<T>) message -> { // code that maps message to T }) .withCoder( // a coder for T ))默认的JmsIO.read()返回PCollection<JmsRecord>,其中包含 JMS headers、properties 以及TextMessage的载荷;而readMessage()则把每条消息交给MessageMapper转成用户自定义类型。MessageMapper<T>继承自Serializable,接口定义在 JmsIO.java 第 632 行附近。
可选的发布重试配置
写入端还提供withRetryConfiguration(RetryConfiguration)用于失败消息重发(源码第 1294-1297 行)。RetryConfiguration默认单次重试间隔 15 秒、最大累计重试时长 1000 天,可通过三种方式创建:
RetryConfiguration retryConfiguration = RetryConfiguration.create(5); // 或 RetryConfiguration retryConfiguration = RetryConfiguration.create(5, Duration.standardSeconds(30), null); // 或 RetryConfiguration retryConfiguration = RetryConfiguration.create(5, Duration.standardSeconds(30), Duration.standardDays(15));参数依次为最大重试次数、单次重试间隔时长、累计重试时长上限。
测试佐证
仓库中 JmsIOTest.java 与集成测试 JmsIOIT.java 覆盖了读写端配置与消息映射逻辑,可作为进一步阅读该连接器行为细节的入口。
新增 Apache PulsarIO(实验性)
2.39.0 通过 BEAM-8218,核心类为 PulsarIO.java。需要说明的是:源码 Javadoc 标注该 IO 目前处于实验性阶段(参见read()、write()的注释,官方跟踪问题为 apache/beam#31078),可能存在 bug 或性能问题,生产环境使用需谨慎评估。
读取端 API
读取端提供两个静态工厂方法:
// 返回 PCollection<PulsarMessage> PulsarIO.read() // 通过 fn 把 Pulsar Message 映射为自定义类型 T PulsarIO.<T>read(fn)关键配置方法(对应源码Read类的 builder):
withClientUrl(String url):Pulsar 客户端地址,例如"pulsar://localhost:6650";withAdminUrl(String url):Admin 客户端地址,例如"http://localhost:8080",可选,用于估算 backlog;withTopic(String topic):读取的 Topic;withStartTimestamp(Long)/withEndTimestamp(Long):按时间范围回溯读取;withEndMessageId(MessageId):按消息 ID 截止读取;withPublishTime():元素时间戳取消息发布时间(默认值);withProcessingTime():元素时间戳取处理时刻;withConsumerPollingTimeout(long):消费者轮询超时,默认 2 秒,调低优化延迟,消费者拉不到数据时可适当调大(源码校验要求大于 0);withPulsarClient(...)/withPulsarAdmin(...):提供自定义的客户端/Admin 工厂函数。
从源码expand()实现(第 202-216 行)可以看出,读取端内部通过Create构造一个PulsarSourceDescriptor,再交给NaiveReadFromPulsarDoFn完成拉取,这也解释了为何它目前属于“朴素实现”的实验性阶段。
写入端 API
PulsarIO.write() .withClientUrl("pulsar://localhost:6650") .withTopic("my-topic");写入端接收PCollection<byte[]>,内部通过WriteToPulsarDoFn逐条发送(源码第 269-273 行)。
测试佐证
PulsarIOTest.java、ReadFromPulsarDoFnTest.java 与集成测试 PulsarIOIT.java 覆盖了读写流程。
BigQueryIO StreamingInserts 线程数优化
2.39.0 通过 BEAM-14283 减少了BigqueryIOStreamingInserts(流式写入 BigQuery)所派生的线程数量,从而降低流式写入场景下的资源开销。这类优化对高频小批量写入 BigQuery 的生产管道尤为有意义,属于运行时资源层面的内部改进,API 用法不受影响。
新特性与改进详解
Flink Scala 2.12 支持
2.39.0 增加了对 Flink Scala 2.12 的支持(BEAM-14386),理由是绝大多数依赖库从 2.12 版本起提供支持。该变更主要服务于使用 Scala 2.12 构建 Flink 作业的 Beam Flink Runner 用户,属于 runner 集成层的兼容性扩展。
Interactive Beam:Dataproc 集群管理的 JupyterLab 扩展
针对 Python SDK 的 Interactive Beam,2.39.0 实现了 JupyterLab 扩展用于“Manage Clusters”,方便用户配置由 Interactive Beam 托管的 Dataproc 集群(BEAM-14130)。该扩展把集群的生命周期管理集成到 Jupyter 工作环境中,降低在 Notebook 里进行大规模交互式探索的运维成本。
Go SDK:Pipeline Drain 支持(实验性)
Go SDK 在 2.39.0 中加入了 Pipeline drain 支持(BEAM-11106)。官方发布说明特别强调:该特性尚未完全验证,本次发布中应视为实验性。Drain 允许管道在停止前先停止接收新输入、等待已有数据处理完成,从而实现优雅下线。
同时,2.39.0 在 Go SDK 的 pipelinex 包中移除了三个无用函数:ShallowCloneParDoPayload()、ShallowCloneSideInput()和ShallowCloneFunctionSpec()(BEAM-13739),属于破坏性变更,升级后引用这些函数的代码需要一并清理。
Python DataFrame API:unstack / pivot
Python SDK 的 DataFrame API 在 2.39.0 中新增了DataFrame.unstack()、DataFrame.pivot()和Series.unstack()(BEAM-13948:
DeferredSeries.unstack()(第 981 行):把 Series 的某一层索引展开为列;DeferredDataFrame.pivot()(第 3812 行):以指定列为 index、columns 进行数据透视,内部通过pivot_helper(第 3926 行)完成分组转换。
这两个方法让用户可以在 Beam 的分布式 DataFrame 上执行与 pandas 语义一致的宽表变换,而无需退回到逐行处理。
Dataflow Runner 模拟凭据支持
Java 与 Python SDK 的 Dataflow Runner 均增加了 impersonation credentials(模拟凭据)支持(BEAM-14014)。该能力允许 Dataflow 作业以服务账号 A 的身份提交、实际以服务账号 B 的权限运行,适用于需要跨项目/跨账号权限隔离的合规场景。
Elasticsearch 8.x 支持
通过 BEAM-14003,ElasticsearchIO 增加了对 Elasticsearch 8.x 的支持,使连接器可以对接较新的 ES 集群版本。
Kinesis 分片感知聚合
Java SDK 的 KinesisIO 实现了 shard-aware 的记录聚合(AWS SDK v2,BEAM-14104),按分片维度聚合写入,减少跨分片请求并提升写入吞吐效率。
ZetaSQL 升级
2.39.0 将 ZetaSQL 升级到 2022.04.1(BEAM-14348),为 Beam SQL 引擎带来更新版本的 ZetaSQL 语义与函数集。
其他修复
- 修复了
ReadFromBigQuery无法与 interactive runner 配合使用的问题(BEAM-14112); - 修复了 Java Spanner IO 在模板执行未指定 ProjectID 时的 NPE(BEAM-14405);
- 修复了
BigQueryServicesImpl.getErrorInfo的潜在 NPE(BEAM-14133)。
破坏性变更与迁移指南
升级到 2.39.0 时,以下破坏性变更需要特别关注。
JmsIO:必须显式配置 valueMapper
2.39.0 起 JmsIO 的写入端强制要求设置valueMapper(BEAM-16308),否则会在管道展开时抛出IllegalArgumentException: withValueMapper() is required。官方迁移示例是使用内置的TextMessageMapper把String输入转成 JMSTextMessage:
JmsIO.<String>write() .withConnectionFactory(jmsConnectionFactory) .withValueMapper(new TextMessageMapper());仓库中 JmsIO.java 的校验逻辑(第 1313 行checkArgument(getValueMapper() != null, "withValueMapper() is required"))与该破坏性变更完全对应。如果你的管道之前只配置了withConnectionFactory+withQueue/withTopic,升级后必须补上valueMapper才能通过校验。
Python SDK:Coder 必须继承 Coder 基类
BEAM-14351 规定 Python SDK 中的 Coder 类需要显式继承Coder基类。这收紧了自定义 Coder 的约束,过去隐式鸭子类型式的 Coder 定义将不再被视为合法 Coder,需要显式class MyCoder(Coder)并实现相应协议。
Python SDK:FileSystem 新增 metadata() 抽象方法
BEAM-14314 在io.filesystem.FileSystem中新增了抽象方法metadata()。任何直接继承FileSystem的自定义文件系统实现都必须实现该方法(通常用于返回文件/目录的元数据信息,如大小、最后修改时间等),否则将因抽象方法未实现而无法实例化。
弃用与版本下线
- Flink 1.11 不再支持(BEAM-14139):Flink Runner 的最低支持版本提升,运行在 Flink 1.11 上的作业需要迁移到受支持的 Flink 版本。
- Python 3.6 不再支持(BEAM-13657):Beam Python SDK 的最低 Python 版本要求提高,仍使用 Python 3.6 的环境需要升级 Python 运行时。
已知问题与完整变更来源
2.39.0 存在一批影响该版本的问题(官方 JIRA 查询条件为project = BEAM AND affectedVersion = 2.39.0,按优先级排序,跟踪问题见 BEAM-14412)。更完整的逐条变更明细可查阅该版本的详细发布说明;需要逐条核对 issue 级别的行为变化时,建议以 JIRA 对应条目的最终状态为准。
升级建议小结
- JMS 用户:为
JmsIO.write()补上valueMapper(如TextMessageMapper);如需按数据路由到不同 Topic,使用withTopicNameMapper并确保与withQueue/withTopic互斥; - Pulsar 用户:可开始试用实验性的
PulsarIO.read()/write(),注意其尚未完全成熟,生产前需充分压测; - Python 用户:检查自定义 Coder 是否显式继承
Coder;若实现了自定义FileSystem,补齐metadata()方法;Python 运行时需高于 3.6; - Flink 用户:确认集群 Flink 版本高于 1.11;使用 Scala 2.12 构建的作业可受益于新增支持;
- Dataflow 用户:可按需启用模拟凭据,实现作业提交身份与运行权限的解耦;
- DataFrame API 用户:
unstack()、pivot()可直接用于宽表化与透视分析场景。
本文涉及的源码入口包括 JmsIO.java、JmsIOTest.java、PulsarIO.java 以及 frames.py,读者可沿这些文件深入 2.39.0 特性的实现细节。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam 中 KafkaIO 连接器:从 Kafka Topic 读取与写入的完整实践指南
Apache Beam 中 KafkaIO 连接器:从 Kafka Topic 读取与写入的完整实践指南 导读 Apache Beam 为 Apache Kaf
大数据批处理流处理数据工程RabbitMQ 3.8.31 维护版本发布解析:升级要点、JMS Topic Exchange 修复与依赖演进
RabbitMQ 3.8.31 维护版本发布解析:升级要点、JMS Topic Exchange 修复与依赖演进 本篇技术指南以 release notes/3
后端消息队列消息路由SeaTunnel RocketMQ 连接器全解析:Source 与 Sink 配置实战、事务写入与版本演进
SeaTunnel RocketMQ 连接器全解析:Source 与 Sink 配置实战、事务写入与版本演进 RocketMQ 是 SeaTunnel Conn
数据集成ETL大数据批处理流处理变更数据捕获
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考