news 2026/10/9 2:28:58

Apache Beam 2.39.0 版本解析:JMS 动态 Topic 写入、PulsarIO 新连接器与 Go SDK 演进

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Beam 2.39.0 版本解析:JMS 动态 Topic 写入、PulsarIO 新连接器与 Go SDK 演进
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

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

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

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 引入两项能力:

  1. 任意类型映射:可以把任何类型的输入元素映射为javax.jms.Message的任意子类(如TextMessage、ObjectMessage、BytesMessage),不再局限于字符串文本;
  2. 动态 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 对应条目的最终状态为准。

升级建议小结

  1. JMS 用户:为JmsIO.write()补上valueMapper(如TextMessageMapper);如需按数据路由到不同 Topic,使用withTopicNameMapper并确保与withQueue/withTopic互斥;
  2. Pulsar 用户:可开始试用实验性的PulsarIO.read()/write(),注意其尚未完全成熟,生产前需充分压测;
  3. Python 用户:检查自定义 Coder 是否显式继承Coder;若实现了自定义FileSystem,补齐metadata()方法;Python 运行时需高于 3.6;
  4. Flink 用户:确认集群 Flink 版本高于 1.11;使用 Scala 2.12 构建的作业可受益于新增支持;
  5. Dataflow 用户:可按需启用模拟凭据,实现作业提交身份与运行权限的解耦;
  6. 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.

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

相关推荐

上一篇:Apereo CAS 命令行 Shell:启动方式、完整命令清单与源码解析
下一篇:5步上手RVC-WebUI:一键部署本地语音克隆

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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/9 2:28:39

CMake FindLibLZMA 模块详解:在项目中集成 LZMA/XZ 压缩库

构建工具开发工具CLI 【免费下载链接】CMake Mirror of CMake upstream repository 项目地址&#xff1a; https://gitcode.com/gh_mirrors/cm/CMake 点击查看 免费下载 导读 本篇文章围绕 CMake 官方仓库中的 FindLibLZMA 查找模块展开&#xff0c;该模块用于在 CMake 构建系…

作者头像 李华
网站建设 2026/10/9 2:27:38

C# TCP粘包拆包 终极满分笔记

一、TCP核心本质&#xff08;必考概念&#xff09;TCP是面向字节流的协议&#xff0c;无消息边界。TCP只保证&#xff1a;数据可靠、有序、不重复。TCP不保证&#xff1a;应用层一次发送多少&#xff0c;接收层就一次读到多少。因此必然产生&#xff1a;粘包、拆包&#xff0c;…

作者头像 李华
网站建设 2026/10/9 2:24:19

基于SSM+Vue的社团管理系统:选题、实现到答辩完整指南

基于SSM Vue的社团管理系统&#xff1a;从选题到答辩的完整干货复盘每年到了毕设季&#xff0c;总有不少同学来问我&#xff1a;“社团管理系统还能做吗&#xff1f;会不会太老套&#xff1f;”我的回答一直是&#xff1a;能做&#xff0c;而且很适合。项目不在于多新奇&#…

作者头像 李华
网站建设 2026/10/9 2:23:19

Java学习进程12

线程游戏的实现 2 关于缓冲区 经过线程游戏的初步设计&#xff0c;直接在窗口分层绘制图像时&#xff0c;画面频繁闪烁&#xff1b;是因为分层绘图按代码顺序逐次刷新&#xff0c;清屏与重绘交替出现&#xff0c;形成视觉频闪&#xff1b;因此引入图像缓冲区&#xff0c;所有图…

作者头像 李华
网站建设 2026/10/9 2:22:29

维度砍一半,检索到底差多少:自测 + 独立信源对账

版权与内容来源声明 本文为原创整理。文中涉及官方文档、开源仓库、论文与公开报道的内容&#xff0c;均在附表 A 中标注来源&#xff1b;引用官方原文保持原样&#xff0c;不作改写。文中命令、版本号与界面截图以本文成文时的实测/核验结果为准&#xff0c;标注「待验证」的部…

作者头像 李华