news 2026/9/23 17:33:52

Apache Pulsar 2.6.0 版本特性深度解析:核心、代理、管理与客户端关键更新

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Pulsar 2.6.0 版本特性深度解析:核心、代理、管理与客户端关键更新
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

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

Apache Pulsar 2.6.0 是社区历经 450+ 次提交打磨后发布的重要里程碑版本,围绕大消息传输、主题级策略、可插拔元数据接口、批量消息精确确认、多监听地址等主题带来了一系列新特性、性能优化与缺陷修复。本文以 2.6.0 发布说明为骨架,结合当前仓库源码与broker.conf配置,逐项梳理核心 Pulsar、Proxy、Admin、Functions、Pulsar SQL 与 Java client 的关键更新,并给出可直接落地的配置与代码示例,帮助你判断升级收益并快速上手新能力。

核心 Pulsar(Core Pulsar)

PIP-37:大消息分块(Large message size support)

PIP-37 通过在生产者端将大消息拆分为多个 chunk(消息块),并在消费端按序拼接还原,使得 Pulsar 能够生产与消费远超单条消息上限的大消息。该特性当前仅支持非共享(non-shared)订阅,且属于客户端侧能力,因此需要将 Pulsar 客户端升级到 2.6.0 及以上版本。

在生产端开启分块:

client.newProducer() .topic("my-topic") .enableChunking(true) .create();

从客户端 API 的注释(ProducerBuilder.java)可以看到官方对该特性的使用建议,这也是 2.6.0 之后长期有效的推荐用法:

  1. 该特性目前仅支持非共享订阅与持久化主题(persistent topic);
  2. 使用分块功能时应关闭 batching,避免两种机制互相干扰;
  3. 客户端在收到 broker 确认前会把已发布消息保留在缓冲区,建议调小maxPendingMessages,避免生产端因缓冲消息占用过大内存;
  4. 为命名空间设置消息 TTL,用于清理因 broker 重启或发布中断而残留的不完整 chunked 消息(也可配置ConsumerBuilder#expireTimeOfIncompleteChunkedMessage);
  5. 消费端同样需要合理配置receiverQueueSizemaxPendingChunkedMessage,以便拼接不完整的大消息。

实现细节见 PR-4440(enableChunking相关实现位于 ProducerBuilderImpl.java)。

PIP-39:命名空间变更事件(System Topic)

在 2.6.0 之前,策略只能设置在命名空间级别,命名空间下所有 topic 都继承该策略;许多用户希望把策略精确到单个 topic。之所以不在 ZooKeeper 上复用命名空间级策略的做法,是为了避免给 ZooKeeper 带来更多负载。

PIP-39 引入 system topic 来存储命名空间变更事件,其初衷是把 topic 策略存到 topic 中而不是 ZooKeeper 里,这也是迈向 topic 级策略(topic level policy)的第一步,后续可以基于该能力轻松扩展 topic 级策略支持。实现细节见 PR-4955。

PIP-45:可插拔元数据接口(Pluggable metadata interface)

为让 Pulsar 摆脱对 ZooKeeper 的强依赖、能够接入其他元数据服务,PIP-45 将ManagedLedger迁移到使用MetadataStore接口,从而打通元数据服务插件化通道。通过MetadataStore接口,可以方便地接入 etcd 等第三方元数据存储。该工作为后续 Pulsar 元数据层(pulsar-metadata模块)的独立演化打下了基础,实现细节见 PR-5358。

PIP-54:批量消息索引级确认

此前 broker 只在批次(batch)消息级别跟踪确认状态:如果某个 batch 中只有部分消息被确认,一旦发生批量消息重投,消费者仍可能再次收到本已确认的消息。PIP-54 支持对 batch 内局部索引(local batch index)进行确认,避免上述重复投递。该特性默认关闭,需要在broker.conf中显式开启:

acknowledgmentAtBatchIndexLevelEnabled=true

当前仓库默认值为false(broker.conf),生产环境按需开启。实现细节见 PR-6052。

PIP-58:消费者自定义消息重试延迟

在线业务处理消息时常出现业务逻辑异常,需要重新消费消息,且希望重试延迟可以灵活控制。此前的通行做法是把消息投递到专门的 retry topic(因为生产端可以指定任意延迟),消费者同时订阅业务 topic 与 retry topic。PIP-58 让消费者可以直接为每条消息设置重试延迟:

Consumer<byte[]> consumer = pulsarClient.newConsumer(Schema.BYTES) .enableRetry(true) .receiverQueueSize(100) .deadLetterPolicy(DeadLetterPolicy.builder() .maxRedeliverCount(maxRedeliveryCount) .retryLetterTopic("persistent://my-property/my-ns/my-subscription-custom-Retry") .build()) .subscribe(); consumer.reconsumeLater(message, 10, TimeUnit.SECONDS);

其中enableRetry(true)deadLetterPolicy(...)均已在 ConsumerBuilder.java 中定义为正式 API,reconsumeLater会将该消息延迟指定时长后重新投递,避免自建 retry topic 的运维负担。实现细节见 PR-6449。

PIP-60:支持 SNI 路由以接入各类代理服务器

此前 Pulsar 不支持使用 Apache Traffic Server(ATS)、HAProxy、Nginx、Envoy 等更具扩展性与安全性的第三方代理。这些代理大多支持 SNI(Server Name Indication)路由,可以在不终结 SSL 连接的情况下把流量路由到目标端。PIP-60 在 Pulsar 客户端加入 SNI 路由支持,使得客户端可以通过 SNI 方式穿透各类 L4/L7 代理直连 broker。实现细节见 PR-6566。

PIP-61:多地址宣告(Advertised multiple addresses)

PIP-61 允许 broker 暴露多个 advertised listener,实现内网/外网流量分离。在broker.conf中配置多个监听地址:

advertisedListeners=internal:pulsar://192.168.1.11:6660,external:pulsar://110.95.234.50:6650

客户端侧通过listenerName指定要使用的监听:

PulsarClient.builder() .serviceUrl(url) .listenerName("internal") .build();

当前仓库的 broker.conf 对该配置补充了说明:值格式必须为<listener_name>:pulsar://<host>:<port>,多个监听以逗号分隔;该配置不能与advertisedAddressbrokerServicePort同时使用;internalListenerName用于指定内部监听名,缺省时 broker 使用第一个监听作为内部监听。实现细节见 PR-6903。

PIP-65:为 Pulsar IO Sources 引入BatchSource

PIP-65 引入BatchSource新接口用于编写基于批次的 connector,同时引入BatchSourceTriggerer接口触发BatchSource的数据采集,并在BatchSourceExecutor中提供系统级实现,为批量型数据源(如定时拉取类 connector)提供了统一的接入模型。实现细节见 PR-7090。

Load balancer:新增ThresholdShedder策略

ThresholdShedder比既有LoadSheddingStrategy更灵活:它先计算集群内各 broker 的平均资源使用率,再把单个 broker 的资源使用率与平均值比较;当某 broker 使用率高于「平均值 + 阈值」时,触发过载卸载(overload shedder)。在broker.conf中启用:

loadBalancerLoadSheddingStrategy=org.apache.pulsar.broker.loadbalance.impl.ThresholdShedder

并可针对性地调整以下参数(当前仓库默认值以 broker.conf 为准):

# The broker resource usage threshold. # When the broker resource usage is greater than the pulsar cluster average resource usage, # the threshold shedder will be triggered to offload bundles from the broker. # It only takes effect in ThresholdShedder strategy. loadBalancerBrokerThresholdShedderPercentage=10 # When calculating new resource usage, the history usage accounts for. # It only takes effect in ThresholdShedder strategy. loadBalancerHistoryResourcePercentage=0.9 # The BandWithIn usage weight when calculating new resource usage. # It only takes effect in ThresholdShedder strategy. loadBalancerBandwithInResourceWeight=1.0 # The BandWithOut usage weight when calculating new resource usage. # It only takes effect in ThresholdShedder strategy. loadBalancerBandwithOutResourceWeight=1.0 # The CPU usage weight when calculating new resource usage. # It only takes effect in ThresholdShedder strategy. loadBalancerCPUResourceWeight=1.0 # The heap memory usage weight when calculating new resource usage. # It only takes effect in ThresholdShedder strategy. loadBalancerMemoryResourceWeight=1.0 # The direct memory usage weight when calculating new resource usage. # It only takes effect in ThresholdShedder strategy. loadBalancerDirectMemoryResourceWeight=1.0 # Bundle unload minimum throughput threshold (MB), avoiding bundle unload frequently. # It only takes effect in ThresholdShedder strategy. loadBalancerBundleUnloadMinThroughputThreshold=10

从源码(ThresholdShedder.java)可以看到该策略的完整算法:先基于loadBalancerHistoryResourcePercentage把历史观测值纳入运行平均值,通过LocalBrokerData#getMaxResourceUsageWithWeight计算各资源加权后的使用率;当某个 broker 的当前/历史使用率超过「平均使用率 +loadBalancerBrokerThresholdShedderPercentage」时,findBundlesForUnloading会提议卸载足够多的 bundle,使该 broker 降到当前平均使用率之下 5%(ADDITIONAL_THRESHOLD_PERCENT_MARGIN = 0.05);同时会跳过近期刚卸载过的 bundle,并受loadBalancerBundleUnloadMinThroughputThreshold的吞吐下限约束,避免频繁卸载。实现细节见 PR-6772。

Key_Shared 订阅:新增一致性哈希分布

此前 Key_Shared 订阅在消费者加入/离开时,通过「切分当前已分配的哈希区间」来重分配键。2.6.0 引入一致性哈希分布:在broker.conf中开启后,键将基于一致性哈希环重新分配给新消费者;默认仍采用自动切分(AUTO_SPLIT)方式。

# On KeyShared subscriptions, with default AUTO_SPLIT mode, use splitting ranges or # consistent hashing to reassign keys to new consumers subscriptionKeySharedUseConsistentHashing=false # On KeyShared subscriptions, number of points in the consistent-hashing ring. # The higher the number, the more equal the assignment of keys to consumers subscriptionKeySharedConsistentHashingReplicaPoints=100

需要说明的是,当前仓库的 broker.conf 中该开关默认值已演进为true,且计划在后续版本中默认启用一致性哈希分布。实现细节见 PR-6791。

Key_Shared 订阅:修复新增消费者时的顺序问题

此前 Key_Shared dispatcher 存在一个顺序性缺陷:当新消费者 c2 加入、旧消费者 c1 离开时,原本分配给 c1 的键消息可能路由到 c2,从而破坏 Key_Shared 订阅的消息顺序投递保证。修复方案是让新消费者以「暂停(paused)」状态加入,直到之前的消息被确认后再开始接收,确保消息按序派发。若你仍希望放宽顺序要求,可在消费端设置:

pulsarClient.newConsumer() .keySharedPolicy(KeySharedPolicy.autoSplitHashRange().setAllowOutOfOrderDelivery(true)) .subscribe();

实现细节见 PR-7106 与 PR-7108。

Key_Shared 订阅:支持键哈希区间读取(Key hash range reading)

该 PR 支持 sticky key hash range reader:broker 只派发「消息键哈希落入指定 key hash range」的消息;且单个 reader 可以指定多个键哈希区间:

pulsarClient.newReader() .topic(topic) .startMessageId(MessageId.earliest) .keyHashRange(Range.of(0, 10000), Range.of(20001, 30000)) .create();

相关客户端配置实现可见 ReaderBuilderImpl.java 与 ReaderConfigurationData.java。实现细节见 PR-5928。

用纯 Java AirCompressor 替代基于 JNI 的压缩库

此前数据压缩依赖 JNI 库,压缩大量小 payload 时 JNI 开销可测且库体积较大。2.6.0 将 LZ4、ZStd、Snappy 的压缩实现替换为 AirCompressor——一个被 Presto 使用的纯 Java 压缩库,降低了压缩路径上的 JNI 开销。实现细节见 PR-5390。

支持多个 Pulsar 集群共享同一个 BookKeeper 集群

该 PR 允许多个 Pulsar 集群通过把 BookKeeper 客户端指向指定 BookKeeper 集群的 ZooKeeper 连接串,来共享同一个 BookKeeper 集群。新增配置bookkeeperMetadataServiceUri用于发现 BookKeeper 集群元数据存储,并用元数据服务 URI 初始化 BookKeeper 客户端:

# Metadata service uri that bookkeeper is used for loading corresponding metadata driver # and resolving its metadata service location. # This value can be fetched using `bookkeeper shell whatisinstanceid` command in BookKeeper cluster. # For example: zk+hierarchical://localhost:2181/ledgers # The metadata service uri list can also be semicolon separated values like below: # zk+hierarchical://zk1:2181;zk2:2181;zk3:2181/ledgers bookkeeperMetadataServiceUri=

该值可通过 BookKeeper 集群中的bookkeeper shell whatisinstanceid命令获取,也支持分号分隔的多地址形式。实现细节见 PR-5985。

支持在订阅追平后删除不活跃主题

此前 Pulsar 只支持删除「没有活跃生产者和订阅」的不活跃主题。该 PR 新增能力:当主题的所有订阅都已追平(caught up)且没有活跃生产者/消费者时,也可以删除不活跃主题,并在broker.conf中暴露删除模式:

# Set the inactive topic delete mode. Default is delete_when_no_subscriptions # 'delete_when_no_subscriptions' mode only delete the topic which has no subscriptions and no active producers # 'delete_when_subscriptions_caught_up' mode only delete the topic that all subscriptions has no backlogs(caught up) # and no active producers/consumers brokerDeleteInactiveTopicsMode=delete_when_no_subscriptions

当前仓库默认仍为delete_when_no_subscriptions(broker.conf)。后续计划支持命名空间级别的该配置。实现细节见 PR-6077。

新增跳过瞬时 OOM 触发的 broker 关闭开关

某个 topic 的高派发速率可能让 broker 瞬时 OOM,这属于瞬态错误,内存释放后几秒内即可恢复。但 2.4(PR-4196)引入的「OOM 时重启 broker」能力在大集群中会造成连锁不稳定:topic 在 broker 间迁移、多个 broker 被重启并波及无关 topic。因此该 PR 提供一个动态开关,跳过 OOM 时的 broker 关闭,避免集群不稳定。实现细节见 PR-6634。

ZooKeeper 缓存过期时间可配置

此前 ZooKeeper 缓存过期时间硬编码,无法按需调整(例如 zk-watch 丢失时希望快速刷新、或希望避免频繁 zk-read、规避 zk 读超时等)。2.6.0 起可在broker.conf中配置:

# ZooKeeper cache expiry time in seconds zooKeeperCacheExpirySeconds=300

当前仓库默认值已演化为-1(broker.conf),表示不自动过期,请按实际缓存刷新需求设置。实现细节见 PR-6668。

批量消息场景下的消费拉取优化

消费者向 broker 发送 fetch 请求时会携带期望的消息条数,但启用 batching 时 broker 在 BookKeeper 或缓存中以 entry 为单位存储数据,消息条数与 entry 数之间存在换算缺口。该 PR 新增avgMessagesPerEntry变量,记录单个 entry 中的平均消息数,在 broker 向消费者推送消息时更新;处理 fetch 请求时据此把请求条数映射为 entry 数,并将avgMessagePerEntry暴露到消费者统计指标 JSON 中。可在broker.conf中开启:

# Precise dispatcher flow control according to history message number of each entry preciseDispatcherFlowControl=false

默认关闭(broker.conf),开启后按历史每 entry 消息数进行精确的 dispatcher 流控。实现细节见 PR-6719。

精确的 topic 发布速率限制

此前 Pulsar 已支持发布速率限制,但控制不够精确。对需要精确限流的场景,可在broker.conf中开启:

preciseTopicPublishRateLimiterEnable=true

默认关闭(broker.conf)。实现细节见 PR-7078。

暴露新 entry 检查延迟

此前新 entry 的检查延迟固定为 10ms 且不可调整。对消费延迟敏感的场景,可在broker.conf中调小(或设为 0):

managedLedgerNewEntriesCheckDelayInMillis=10

注意:取值越小,消费吞吐可能越低,需要权衡。实现细节见 PR-7154。

Schema:KeyValue schema 支持null键与null

该 PR 让 KeyValue schema 支持键或值为 null 的消息,补全了 KeyValue 场景的数据建模能力。实现细节见 PR-7139。

支持到达maxLedgerRolloverTimeMinutes时触发 ledger 翻转

该 PR 实现一个监控线程,周期性检查当前 topic ledger 是否满足managedLedgerMaxLedgerRolloverTimeMinutes约束并触发翻转,使配置真正生效。其核心收益在于:翻转后可以关闭当前 ledger,从而释放当前 ledger 的存储空间。对低频 topic 而言,其当前 ledger 数据很可能早已过期,而旧逻辑只在追加新 entry 时触发翻转,明显浪费磁盘。监控线程按固定时间间隔调度,间隔即为managedLedgerMaxLedgerRolloverTimeMinutes;每次检查同时做两个判断:currentLedgerEntries > 0currentLedgerIsFull()。当当前 entry 数为 0 时不触发新翻转,可用于减少 ledger 创建。

当前仓库配置项位于 broker.conf,默认值为240(分钟);ledger 翻转相关接线可在 BrokerService.java 中看到(setMaximumRolloverTimesetMinimumRolloverTimesetLedgerRolloverTimeout等)。实现细节见 PR-7111/PR-7116。

Proxy

新增获取连接与 topic 统计的 REST API

此前 Pulsar proxy 缺少获取自身内部信息的统计接口。该 PR 为 proxy 新增 REST API,可获取实时连接数、topic 统计(更高日志级别下)等信息,便于运维观察代理层状态。实现细节见 PR-6473。

Admin

pulsar-admin 支持按消息 ID 获取消息

该 PR 为 pulsar-admin 新增get-message-by-id命令,用户提供 ledger ID 与 entry ID 即可查看单条消息。命令注册位于 CmdTopics.java,可用于排查特定 entry 的存储内容。实现细节见 PR-6331。

支持强制删除订阅

该 PR 新增deleteForcefully方法,支持强制删除订阅(即使订阅仍存在积压)。实现细节见 PR-6383。

Functions

  • 内置函数(Built-in functions):以与内置 connector 相同的方式创建内置函数,简化函数的分发与加载。实现细节见 PR-6895。
  • Go Function 心跳(gRPC 服务):为生产环境使用增加 Go Function 心跳与 gRPC 服务,提升 Go 运行时函数实例的可观测性与管理能力。实现细节见 PR-6031。
  • 函数自定义属性选项:提交函数时允许设置自定义系统属性,可用于通过系统属性传递凭据。实现细节见 PR-6348。
  • 函数 worker 与 broker 的 TLS 配置分离:将函数 worker 与 broker 的 TLS 配置拆开,便于独立管理两者安全配置。实现细节见 PR-6602。
  • 函数与 source 中构建消费者:此前函数与 source 的 context 只允许创建 publisher 而不允许创建 consumer,该 PR 补齐此能力。实现细节见 PR-6954。

Pulsar SQL

支持 KeyValue schema

此前 Pulsar SQL 无法读取 KeyValue schema 数据。该 PR 为 Pulsar SQL 增加 KeyValue schema 支持:key 字段名加key.前缀,value 字段名加value.前缀,从而可以在 SQL 中区分键与值。实现细节见 PR-6325。

支持多个 Avro schema 版本

此前如果 topic 存在多个 Avro schema 版本,用 Pulsar SQL 查询该 topic 会引入问题。该改动后,可以演化 topic 的 schema,并在查询时保持该 topic 所有 schema 的传递后向兼容。实现细节见 PR-4847。

Java client

关闭 producer 时支持等待 in-flight 消息

此前关闭 producer 时,pulsar-client 会立即失败所有 in-flight 消息——即使这些消息已在 broker 端持久化成功。多数场景下用户更希望等待这些消息完成而不是失败。该 PR 为 close API 增加标志位,控制关闭时是否等待 in-flight 消息:开启后关闭 producer 会等待 in-flight 消息,pulsar-client 不会立刻判定这些消息失败。实现细节见 PR-6648。

支持从输入流动态加载 TLS 证书/密钥

此前默认 TLS 认证 providerAuthenticationTls只接受证书与密钥的文件路径,但某些应用难以在本地保存证书/密钥文件。该 PR 为AuthenticationTls增加流(stream)支持,可提供 X509Certs 与 PrivateKey,并在给定 provider 的流内容变化时自动刷新。实现细节见 PR-6760。

异步发送异常时返回 sequence ID

此前异步发送失败抛出的异常无法定位是哪条消息异常,用户难以判断需要重试哪些消息。该 PR 在客户端侧作出改进:抛异常时在org.apache.pulsar.client.api.PulsarClientException中设置sequenceId,便于精准定位与重试。实现细节见 PR-6825。

升级与获取

  • 下载 Apache Pulsar 2.6.0 可通过 Apache Pulsar 官网下载页面获取。
  • 更多细节可查阅 Apache Pulsar 2.6.0 release notes 与 milestone 为 2.6.0 的 PR 列表。
  • 升级前请重点评估本版本中的默认行为变化:例如 Key_Shared 一致性哈希开关、ZooKeeper 缓存过期时间默认值、preciseDispatcherFlowControl等,均可在 conf/broker.conf 中按上文说明调整。

如有问题或建议,可通过users@pulsar.apache.orgdev@pulsar.apache.org邮件列表或 Pulsar Slack 社区与项目组交流,并欢迎为 Pulsar 贡献代码。

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

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载
上一篇:DotNetGuide代码生成技术深度解析
下一篇:APatch完全指南:Android内核级Root终极解决方案详解

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

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

ABAQUS模拟钢制重力锚在钙质土中的承载力分析

1. 项目概述在深海工程领域&#xff0c;重力锚作为固定海底管道、电缆和浮式结构的关键部件&#xff0c;其承载性能直接关系到整个工程系统的安全性和可靠性。钙质土作为一种特殊的海洋沉积物&#xff0c;广泛分布于热带和亚热带海域&#xff0c;其力学特性与常规陆相土体存在显…

作者头像 李华
网站建设 2026/9/23 17:30:48

OpenCV人脸识别考勤系统实战:从环境搭建到落地避坑

简介&#xff1a;这份资源是面向高校学生与Python初学者的人脸识别考勤系统完整项目源码&#xff0c;适合用作课程设计、期末大作业或OpenCV与dlib入门实战参考。项目围绕考勤管理场景&#xff0c;实现了用户注册登录、人脸检测与识别、打卡记录及数据查询等核心功能&#xff0…

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

双目立体视觉毕设指南:标定、匹配与深度图生成

简介&#xff1a;这份资源是面向计算机、人工智能、自动化、电子信息等专业学生与科研人员的双目摄像头立体视觉系统完整项目包&#xff0c;围绕相机标定、立体匹配与深度图生成三大核心环节展开&#xff0c;可作为毕业设计、课程设计或项目立项演示的参考方案。压缩包共190个文…

作者头像 李华
网站建设 2026/9/23 17:20:00

ID3DXSprite实战:Visual C++打造实时DSP数据可视化显示

简介&#xff1a;针对Direct3D中ID3DXSprite接口的2D图形渲染与DirectSound编程实践&#xff0c;这份Visual C示例工程适用于游戏开发和实时界面设计场景&#xff0c;帮助开发者解决精灵批量绘制效率与音频处理问题。包内共43个文件&#xff0c;以bmp位图资源、cpp源文件、h头文…

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

DeepSeek本地部署与API调用实战:从模型选型到避坑指南

简介&#xff1a;这份PDF由清华大学新闻与传播学院新媒体研究中心整理&#xff0c;以DeepSeek-R1开源推理模型为主线&#xff0c;系统展示智能对话、文本生成、代码补全、知识推理、联网搜索与文件读取等能力&#xff0c;并专门对比了推理模型与非推理模型在快慢思考、创造力、…

作者头像 李华