- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
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;该参数的类型为SubscriptionInitialPosition(LATEST/EARLIEST)。在创建 Sink 配置时,该值会被写入SinkConfig的sourceSubscriptionPosition字段(见 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),同步与异步形式分别为getAllSchemas与getAllSchemasAsync;testCompatibility:在不注册 schema 的前提下测试其兼容性。仓库接口定义见 Schemas.java 与 Schemas.java,返回类型为IsCompatibilityResponse,支持PostSchemaPayload与SchemaInfo两种入参形式;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 中补齐了返回MessageID的send()接口。
对于正在评估是否升级到 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
相关推荐
Apache Pulsar Schema 机制深度解析:SchemaInfo、Schema 类型与版本演进原理
Apache Pulsar Schema 机制深度解析:SchemaInfo、Schema 类型与版本演进原理 本指南基于 Apache Pulsar 官方文档
消息队列后端流处理Apache Pulsar 2.8.1 版本全解析:Broker、Proxy、Functions 与客户端的关键修复与增强
Apache Pulsar 2.8.1 版本全解析:Broker、Proxy、Functions 与客户端的关键修复与增强 本篇技术指南以 Apache Pul
消息队列后端流处理Apache Pulsar 2.6.1 版本深度解析:Broker、客户端与 Functions 的关键修复与增强
Apache Pulsar 2.6.1 版本深度解析:Broker、客户端与 Functions 的关键修复与增强 Apache Pulsar 2.6.1 是社
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考