news 2026/9/23 14:58:47

Apache Pulsar 2.4.2 版本要点解析:Functions 类加载、去重修复、Schema 新 API 与消费位置控制

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Pulsar 2.4.2 版本要点解析:Functions 类加载、去重修复、Schema 新 API 与消费位置控制
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

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

Apache Pulsar 2.4.2 是 Pulsar 社区在 2.4.x 系列中的一个重要维护版本,包含 110 余次提交,聚焦于 Functions 运行时、消息去重(Deduplication)、Schema 管理 API、Broker 与订阅管理等多处改进与缺陷修复。本文以该版本的官方发布博客为核心脉络,结合当前仓库中的源码、客户端 API 与 CLI 实现,逐条解读这些改进背后的设计动机、使用方式与实现依据,帮助读者理解 2.4.2 版本在运行时可靠性与开发体验上的具体变化。

使用 classLoaders 加载 Java Functions

在 Pulsar 2.4.2 中,无论 Java Functions 实例使用 shaded JAR 还是 classLoader 方式启动,窗口函数(windowed functions)都能正常工作;同时,当启用--output-serde-classname选项时,functionClassLoader也会被正确设置。

在 2.4.2 之前,Java Functions 实例默认以 shaded JAR 启动,运行时会使用不同的 classLoader 分别加载 Pulsar 内部代码、用户代码以及两者交互所需的接口,由此产生两个问题:

  • 当 Java Functions 实例改用 classLoader 方式时,窗口函数无法正常工作;
  • 使用--output-serde-classname选项时,functionClassLoader没有被正确设置,导致输出序列化/反序列化器无法从正确的类加载上下文中加载。

该修复使得 Functions 的运行时隔离与序列化链路在两种启动方式下行为一致。窗口函数依赖时间窗口与计数窗口的批处理语义,其正确性直接取决于用户代码、Pulsar API 与运行时实现是否处于同一可预期的类加载层级中。

启用 TLS 后以 Functions worker 启动 Broker

2.4.2 支持在 Broker 客户端启用 TLS 的场景下,将 Broker 与 Functions worker 一同启动。

此前,当 Functions worker 与 Broker 一起运行时,worker 会在function_worker.yml文件中检查是否启用了 TLS;若启用则使用 TLS 端口。但当 Functions worker 自身启用 TLS 时,它检查的却是broker.conf。由于此时 Functions worker 实际是随 Broker 一起运行的,以broker.conf作为判断是否使用 TLS 的单一事实来源(single source of truth)才是合理的——这正是 2.4.2 的修复方式:统一从broker.conf读取 TLS 配置,避免配置来源不一致导致端口或协议判断错误。

相关配置可参考仓库中的 conf/broker.conf 与 conf/functions_worker.yml。

读取不存在的函数状态 Key 时返回错误码与错误信息

Pulsar Functions 支持使用 BookKeeper 持久化函数状态(function state)。此前,当用户尝试从函数状态中读取一个不存在的 Key 时,会直接抛出 NPE(NullPointerException),错误信息不明确。

2.4.2 为“Key 不存在”这一场景补充了明确的错误码与错误信息,使上层应用能够区分“状态不存在”与“真实异常”,从而进行更合理的容错处理,而不是依赖对 NPE 的捕获与猜测。

消息去重(Deduplication)的可靠性修复

Pulsar 的消息去重机制基于“已持久化的最大 Sequence ID”来过滤重复消息:只有 Sequence ID 大于已持久化最大值的消息才会被接受。但该机制在出现错误时存在一个隐患——如果某条消息的持久化过程失败,但失败信息本身被误当作“已去重”,那么生产者的重试消息会被“去重”掉,导致该消息永远无法真正落盘。

2.4.2 从两个方面修复了这一问题:

  • 双重校验待处理消息:当去重状态不确定时(例如消息仍处于 pending 状态),会向生产者返回错误,而不是静默丢弃;
  • 失败后同步 lastPushed 与 lastStored 映射:持久化失败后,将 lastPushed 映射与 lastStored 映射对齐,避免游标位置不一致导致后续消息被错误去重。

该逻辑位于 Broker 的消息生产链路中。修复后,去重在异常场景下优先“报错”而非“丢消息”,兼顾了去重的性能收益与消息投递的可靠性。

Sink 支持从最早位置消费数据

2.4.2 为 Pulsar Sinks 新增了--subs-position参数,使用户可以指定从最新(latest)或最早(earliest)位置消费数据。

在此之前,Sink 默认从 topic 的最新位置开始消费,无法消费 Sink topic 中更早的历史数据。在 2.4.2 中,Sink CLI 提供了该参数。以当前仓库中 pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSinks.java 为例,其定义如下:

@Parameter(names = "--subs-position", description = "Pulsar source subscription position " + "if user wants to consume messages from the specified location") protected SubscriptionInitialPosition subsPosition;

该参数的类型为SubscriptionInitialPositionLATEST/EARLIEST)。在创建 Sink 配置时,该值会被写入SinkConfigsourceSubscriptionPosition字段(见 CmdSinks.java):

if (null != subsPosition) { sinkConfig.setSourceSubscriptionPosition(subsPosition); }

对应的 Functions 运行时配置中也有subscriptionPosition字段(见 pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/source/PulsarSourceConfig.java),默认行为为LATEST

使用示例:

bin/pulsar-admin sinks create \ --tenant public \ --namespace default \ --name my-sink \ --sink-type <type> \ --inputs persistent://public/default/input-topic \ --subs-name my-sink-subscription \ --subs-position earliest

需要说明的是,--subs-position仅在订阅不存在时生效;如果指定的订阅已经存在,消费位置将沿用该订阅已有的游标。

订阅类型变更时关闭旧的 Dispatcher

在 2.4.2 中,当 topic 的订阅类型发生变更时,Broker 会创建新的 dispatcher 并关闭旧 dispatcher,从而避免内存泄漏。

此前,订阅类型变更时新 dispatcher 被创建、旧 dispatcher 被直接丢弃而未被关闭。如果订阅游标不是持久化的(non-durable),当所有消费者移除后,该订阅会从 topic 上关闭并移除,此时 dispatcher 也应被关闭;否则其内部的 RateLimiter 实例不会被垃圾回收,最终造成内存泄漏。2.4.2 通过在订阅变更及订阅移除路径上显式关闭旧 dispatcher 修复了该问题。

基于订阅顺序选择活跃消费者(Active Consumer)

Pulsar 的 Failover 订阅模式要求从消费者列表中选举一个“活跃消费者”(active consumer)来接收消息。

在 2.4.2 之前,活跃消费者是基于优先级(priority level)与消费者名称排序后选出的。在这种方式下,活跃消费者加入、离开时,可能出现没有任何消费者被真正选举为“活跃”或消费到消息的异常状态。

2.4.2 改为直接按订阅顺序选择活跃消费者:即消费者列表中的第一个消费者被选为活跃消费者,无需额外排序。这使得 Failover 语义更直观、稳定。

从连接中正确移除失败的 Producer

2.4.2 修复了 Broker 无法从连接中正确清理旧失败 Producer 的问题。

此前,当 Broker 尝试清理失败 Producer 中的producer-future时,会错误地移除新创建的producer-future,而不是旧的失败 Producer,导致 Broker 日志中出现如下告警:

17:22:00.700 [pulsar-io-21-26] WARN org.apache.pulsar.broker.service.ServerCnx - [/1.1.1.1:1111][453] Producer with id persistent://prop/cluster/ns/topic is already present on the connection

当前仓库的 Broker 实现中保留了类似的“Producer/Consumer already present on the connection”告警路径(见 pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java),这正是该问题所涉及的生产连接注册/去注册逻辑。2.4.2 修复后,失败的 Producer 会被正确地从连接中移除,新 Producer 注册时不会再与残留条目冲突。

Schema 管理新增三个 API

2.4.2 在 Schema 管理方面新增了三个 Admin API,全部可以在 pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Schemas.java 中看到对应定义:

  • getAllVersions:返回给定 topic 的 schema 版本列表。仓库中对应的接口为getAllSchemas(String topic)(见 Schemas.java),同步与异步形式分别为getAllSchemasgetAllSchemasAsync
  • testCompatibility:在不注册 schema 的前提下测试其兼容性。仓库接口定义见 Schemas.java 与 Schemas.java,返回类型为IsCompatibilityResponse,支持PostSchemaPayloadSchemaInfo两种入参形式;
  • getVersionBySchema:给定一个 schema 定义,返回其对应的 schema 版本。仓库接口定义见 Schemas.java 与 Schemas.java。

三个 API 均提供同步与异步(Async后缀)两种形式,返回CompletableFuture,便于在异步场景下使用。这些 API 的意义在于:testCompatibility让开发者在正式注册 schema 前即可验证兼容性策略;getVersionBySchema支持反查某个 schema 定义对应的版本号;getAllVersions则便于审计 topic 上 schema 的演进历史。

Consumer 暴露getLastMessageId()方法

2.4.2 在ConsumerImpl中暴露了getLastMessageId()方法。该方法在客户端 API 层已有正式定义(见 pulsar-client-api/src/main/java/org/apache/pulsar/client/api/Consumer.java):

MessageId getLastMessageId() throws PulsarClientException; CompletableFuture<MessageId> getLastMessageIdAsync();

对用户而言,该方法的实用价值在于:

  • 计算消息积压(lag):将getLastMessageId()与当前已消费的位置比较,即可获知滞后消息数量;
  • 只消费当前时间之前的消息:以当前最后一条消息 ID 为边界,控制消费范围。

例如,可以配合 Reader 或 Consumer 在启动时获取 topic 的最新消息 ID,从而决定是否跳过后续新到达的消息。

C++/Go 客户端新增send()接口返回 MessageID

2.4.2 在 C++ 与 Go 客户端中新增了send()接口,使发送成功后能够将MessageID返回给用户,其行为与 Java 客户端保持一致——在 Java 中,MessageId send(byte[] message)会将MessageId返回给调用者。

这一改动统一了多语言客户端的发送语义:此前 C++/Go 的发送接口只返回发送结果,用户无法直接拿到消息 ID 用于后续的追踪、幂等或游标记录;2.4.2 之后,三个主流客户端的行为对齐,降低了跨语言迁移时的认知成本。

订阅失败后取消 Consumer 后台任务

在 2.4.2 中,确保订阅失败后 Consumer 的后台任务被正确取消。

此前,ConsumerImpl构造函数会启动一些后台任务;但如果消费者创建失败,这些任务并不会被取消,导致对象上仍保留活跃引用,造成资源泄漏。2.4.2 在消费者创建/订阅失败路径上补充了任务取消逻辑,保证失败场景下后台任务与相关对象能被及时清理。

删除附着 regex 消费者的 Topic

2.4.2 支持删除附着有正则消费者(regex consumer)的 topic。此前这是不可能做到的,原因在于:regex consumer 在 topic 被删除后会立即重新连接并重新创建该 topic。

实现上,2.4.2 通过两个手段解决:

  • CommandSubscribe中增加一个标志位,使 regex consumer 永远不会触发 topic 的自动创建;
  • 订阅一个不存在的 topic 时,当特定错误发生时,将该消费者判定为永久失败并停止重试。

这样,删除 topic 后 regex consumer 不会“复活”它,topic 才能真正被清理。从源码结构看,该能力涉及客户端订阅协议(CommandSubscribe)与 Broker 端 topic 自动创建策略的协同修改。

小结

Apache Pulsar 2.4.2 是一轮以“运行时可靠性与开发体验”为主题的版本更新:

  • Functions 侧:统一了 classLoader/shaded JAR 两种启动方式下的类加载行为,修正了函数状态 Key 缺失时的错误反馈;
  • Broker 侧:修复了去重异常场景下可能丢消息的隐患、订阅类型变更时的 dispatcher 内存泄漏、失败 Producer 清理错乱以及活跃消费者选举不稳定等问题;
  • 客户端与 API 侧:新增 Schema 的三个管理 API、Consumer 的getLastMessageId()、Sink 的--subs-position,并在 C++/Go 中补齐了返回MessageIDsend()接口。

对于正在评估是否升级到 2.4.2 的用户,本文涉及的几项修复(尤其去重可靠性、dispatcher 内存泄漏、regex 消费者阻塞 topic 删除)均属于生产环境可能直接遇到的问题;若你正在使用这些特性,建议重点回归验证对应场景。版本发布说明可参见仓库内 site2/website-next/blog/2019-12-04-Apache-Pulsar-2-4-2.md 及官方发布说明。

  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

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

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

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

脉冲S参数测量技术:原理、应用与实战解析

1. 脉冲S参数测量的核心价值与挑战在射频微波领域&#xff0c;脉冲S参数测量正成为评估高速数字电路、雷达系统和功率放大器动态特性的重要手段。与传统连续波测量相比&#xff0c;脉冲激励能更真实地模拟器件在实际工作中的瞬态响应。我曾在某相控阵雷达T/R组件测试中深有体会…

作者头像 李华
网站建设 2026/9/23 14:56:28

化疗药物激活STING通路重塑肿瘤免疫微环境

1. 研究背景与临床意义化疗药物与免疫系统的相互作用一直是肿瘤治疗领域的热点课题。吉西他滨和顺铂作为临床常用的标准化疗方案&#xff0c;在多种实体瘤&#xff08;如非小细胞肺癌、膀胱癌等&#xff09;中展现出明确的治疗效果。但传统观点认为&#xff0c;化疗主要通过直接…

作者头像 李华
网站建设 2026/9/23 14:54:25

农学论文真相✅别堆砌田间方案+作物生理数据硬凑综述

农学、作物栽培学、育种学、植物营养、耕作学方向本科生研究生狠狠共情&#xff01; 农学文献综述&#xff0c;是农林类典型“栽培方案同质化、试验描述高度撞文”的查重重灾区&#xff01; 综述高频覆盖&#xff1a;作物栽培调控、品种选育、施肥管理、抗旱抗逆生理、耕作模式…

作者头像 李华
网站建设 2026/9/23 14:54:23

JSP+SSH+MVC商城源码解析:从分层架构到部署避坑指南

简介&#xff1a;这是一个面向Java Web初学者的水果销售商城系统完整源码包&#xff0c;基于SSH框架与MVC分层设计&#xff0c;涵盖普通用户注册登录、商品分类浏览、购物车下单、订单查询及管理员端的水果增删改查、分类/订单/用户管理等核心业务模块&#xff0c;很适合用做课…

作者头像 李华
网站建设 2026/9/23 14:53:49

Python本地财务管理系统:SQLite+规则引擎实现个人数据自主

简介&#xff1a;这是一套面向计算机专业本科生与毕业设计初学者的Python个人财务管理系统实战项目&#xff0c;聚焦财务管理场景下的工程化开发实践&#xff0c;融合基础业务逻辑与轻量级AI应用思路。资源共30个文件&#xff0c;包含10个核心Python模块&#xff08;如账单处理…

作者头像 李华