
大数据批处理流处理数据工程【免费下载链接】beamApache 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/OsJMS 连接器能力大幅增强、新增 Apache PulsarIO、BigQueryIO 流式写入线程优化New Features / ImprovementsFlink Scala 2.12 支持、Go SDK Pipeline drain、Python DataFrame API 扩展、Dataflow 模拟凭据、Elasticsearch 8.x 支持、Kinesis 分片感知聚合等Breaking ChangesJmsIO 强制要求valueMapper、Python Coder 继承约束、FileSystem 新增metadata()抽象方法、Go pipelinex 无用函数移除DeprecationsFlink 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类型为SerializableFunctionEventT, String接收一个输入事件对象返回该事件应该发往的 Topic 名称。它必须与withQueue()、withTopic()三者互斥使用。valueMapper类型为SerializableBiFunctionEventT, 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 的典型用法SerializableFunctionCompanyEvent, String topicNameMapper (event - String.format( company/%s/employee/%s, event.getCompanyName(), event.getEmployeeId())); pipeline .apply(...) // PCollectionCompanyEvent .apply(JmsIO.write() .withConnectionFactory(jmsConnectionFactory) .withTopicNameMapper(topicNameMapper) .withValueMapper(valueMapper));其中valueMapper可以把事件序列化为 JMSTextMessageSerializableBiFunctionSomeEventObject, 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(MessageMapperT)完成转换源码 JmsIO.java 第 106-120 行的 Javadoc 示例pipeline.apply(JmsIO.TreadMessage() .withConnectionFactory(myConnectionFactory) .withQueue(my-queue) .withMessageMapper((MessageMapperT) message - { // code that maps message to T }) .withCoder( // a coder for T ))默认的JmsIO.read()返回PCollectionJmsRecord其中包含 JMS headers、properties 以及TextMessage的载荷而readMessage()则把每条消息交给MessageMapper转成用户自定义类型。MessageMapperT继承自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读取端提供两个静态工厂方法// 返回 PCollectionPulsarMessage PulsarIO.read() // 通过 fn 把 Pulsar Message 映射为自定义类型 T PulsarIO.Tread(fn)关键配置方法对应源码Read类的 builderwithClientUrl(String url)Pulsar 客户端地址例如pulsar://localhost:6650withAdminUrl(String url)Admin 客户端地址例如http://localhost:8080可选用于估算 backlogwithTopic(String topic)读取的 TopicwithStartTimestamp(Long)/withEndTimestamp(Long)按时间范围回溯读取withEndMessageId(MessageId)按消息 ID 截止读取withPublishTime()元素时间戳取消息发布时间默认值withProcessingTime()元素时间戳取处理时刻withConsumerPollingTimeout(long)消费者轮询超时默认 2 秒调低优化延迟消费者拉不到数据时可适当调大源码校验要求大于 0withPulsarClient(...)/withPulsarAdmin(...)提供自定义的客户端/Admin 工厂函数。从源码expand()实现第 202-216 行可以看出读取端内部通过Create构造一个PulsarSourceDescriptor再交给NaiveReadFromPulsarDoFn完成拉取这也解释了为何它目前属于“朴素实现”的实验性阶段。写入端 APIPulsarIO.write() .withClientUrl(pulsar://localhost:6650) .withTopic(my-topic);写入端接收PCollectionbyte[]内部通过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 BeamDataproc 集群管理的 JupyterLab 扩展针对 Python SDK 的 Interactive Beam2.39.0 实现了 JupyterLab 扩展用于“Manage Clusters”方便用户配置由 Interactive Beam 托管的 Dataproc 集群BEAM-14130。该扩展把集群的生命周期管理集成到 Jupyter 工作环境中降低在 Notebook 里进行大规模交互式探索的运维成本。Go SDKPipeline Drain 支持实验性Go SDK 在 2.39.0 中加入了 Pipeline drain 支持BEAM-11106。官方发布说明特别强调该特性尚未完全验证本次发布中应视为实验性。Drain 允许管道在停止前先停止接收新输入、等待已有数据处理完成从而实现优雅下线。同时2.39.0 在 Go SDK 的 pipelinex 包中移除了三个无用函数ShallowCloneParDoPayload()、ShallowCloneSideInput()和ShallowCloneFunctionSpec()BEAM-13739属于破坏性变更升级后引用这些函数的代码需要一并清理。Python DataFrame APIunstack / pivotPython SDK 的 DataFrame API 在 2.39.0 中新增了DataFrame.unstack()、DataFrame.pivot()和Series.unstack()BEAM-13948DeferredSeries.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-14003ElasticsearchIO 增加了对 Elasticsearch 8.x 的支持使连接器可以对接较新的 ES 集群版本。Kinesis 分片感知聚合Java SDK 的 KinesisIO 实现了 shard-aware 的记录聚合AWS SDK v2BEAM-14104按分片维度聚合写入减少跨分片请求并提升写入吞吐效率。ZetaSQL 升级2.39.0 将 ZetaSQL 升级到 2022.04.1BEAM-14348为 Beam SQL 引擎带来更新版本的 ZetaSQL 语义与函数集。其他修复修复了ReadFromBigQuery无法与 interactive runner 配合使用的问题BEAM-14112修复了 Java Spanner IO 在模板执行未指定 ProjectID 时的 NPEBEAM-14405修复了BigQueryServicesImpl.getErrorInfo的潜在 NPEBEAM-14133。破坏性变更与迁移指南升级到 2.39.0 时以下破坏性变更需要特别关注。JmsIO必须显式配置 valueMapper2.39.0 起 JmsIO 的写入端强制要求设置valueMapperBEAM-16308否则会在管道展开时抛出IllegalArgumentException: withValueMapper() is required。官方迁移示例是使用内置的TextMessageMapper把String输入转成 JMSTextMessageJmsIO.Stringwrite() .withConnectionFactory(jmsConnectionFactory) .withValueMapper(new TextMessageMapper());仓库中 JmsIO.java 的校验逻辑第 1313 行checkArgument(getValueMapper() ! null, withValueMapper() is required)与该破坏性变更完全对应。如果你的管道之前只配置了withConnectionFactorywithQueue/withTopic升级后必须补上valueMapper才能通过校验。Python SDKCoder 必须继承 Coder 基类BEAM-14351 规定 Python SDK 中的 Coder 类需要显式继承Coder基类。这收紧了自定义 Coder 的约束过去隐式鸭子类型式的 Coder 定义将不再被视为合法 Coder需要显式class MyCoder(Coder)并实现相应协议。Python SDKFileSystem 新增 metadata() 抽象方法BEAM-14314 在io.filesystem.FileSystem中新增了抽象方法metadata()。任何直接继承FileSystem的自定义文件系统实现都必须实现该方法通常用于返回文件/目录的元数据信息如大小、最后修改时间等否则将因抽象方法未实现而无法实例化。弃用与版本下线Flink 1.11 不再支持BEAM-14139Flink Runner 的最低支持版本提升运行在 Flink 1.11 上的作业需要迁移到受支持的 Flink 版本。Python 3.6 不再支持BEAM-13657Beam 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.6Flink 用户确认集群 Flink 版本高于 1.11使用 Scala 2.12 构建的作业可受益于新增支持Dataflow 用户可按需启用模拟凭据实现作业提交身份与运行权限的解耦DataFrame API 用户unstack()、pivot()可直接用于宽表化与透视分析场景。本文涉及的源码入口包括 JmsIO.java、JmsIOTest.java、PulsarIO.java 以及 frames.py读者可沿这些文件深入 2.39.0 特性的实现细节。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐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大数据批处理流处理变更数据捕获上一篇Apereo CAS 命令行 Shell启动方式、完整命令清单与源码解析下一篇5步上手RVC-WebUI一键部署本地语音克隆创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考