消息队列后端流处理【免费下载链接】pulsarApache 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 启动 Broker2.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 时会直接抛出 NPENullPointerException错误信息不明确。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 配置时该值会被写入SinkConfig的sourceSubscriptionPosition字段见 CmdSinks.javaif (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 ConsumerPulsar 的 Failover 订阅模式要求从消费者列表中选举一个“活跃消费者”active consumer来接收消息。在 2.4.2 之前活跃消费者是基于优先级priority level与消费者名称排序后选出的。在这种方式下活跃消费者加入、离开时可能出现没有任何消费者被真正选举为“活跃”或消费到消息的异常状态。2.4.2 改为直接按订阅顺序选择活跃消费者即消费者列表中的第一个消费者被选为活跃消费者无需额外排序。这使得 Failover 语义更直观、稳定。从连接中正确移除失败的 Producer2.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 管理新增三个 API2.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与getAllSchemasAsynctestCompatibility在不注册 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.javaMessageId getLastMessageId() throws PulsarClientException; CompletableFutureMessageId getLastMessageIdAsync();对用户而言该方法的实用价值在于计算消息积压lag将getLastMessageId()与当前已消费的位置比较即可获知滞后消息数量只消费当前时间之前的消息以当前最后一条消息 ID 为边界控制消费范围。例如可以配合 Reader 或 Consumer 在启动时获取 topic 的最新消息 ID从而决定是否跳过后续新到达的消息。C/Go 客户端新增send()接口返回 MessageID2.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 消费者的 Topic2.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 及官方发布说明。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐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),仅供参考
