- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
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 之后长期有效的推荐用法:
- 该特性目前仅支持非共享订阅与持久化主题(persistent topic);
- 使用分块功能时应关闭 batching,避免两种机制互相干扰;
- 客户端在收到 broker 确认前会把已发布消息保留在缓冲区,建议调小
maxPendingMessages,避免生产端因缓冲消息占用过大内存; - 为命名空间设置消息 TTL,用于清理因 broker 重启或发布中断而残留的不完整 chunked 消息(也可配置
ConsumerBuilder#expireTimeOfIncompleteChunkedMessage); - 消费端同样需要合理配置
receiverQueueSize与maxPendingChunkedMessage,以便拼接不完整的大消息。
实现细节见 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>,多个监听以逗号分隔;该配置不能与advertisedAddress、brokerServicePort同时使用;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 > 0与currentLedgerIsFull()。当当前 entry 数为 0 时不触发新翻转,可用于减少 ledger 创建。
当前仓库配置项位于 broker.conf,默认值为240(分钟);ledger 翻转相关接线可在 BrokerService.java 中看到(setMaximumRolloverTime、setMinimumRolloverTime、setLedgerRolloverTimeout等)。实现细节见 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.org、dev@pulsar.apache.org邮件列表或 Pulsar Slack 社区与项目组交流,并欢迎为 Pulsar 贡献代码。
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Roo Code attempt_completion 工具全解析:任务收尾、结果呈现与迭代反馈机制
Roo Code attempt_completion 工具全解析:任务收尾、结果呈现与迭代反馈机制 attempt_completion 是 Roo Code
消息队列后端流处理解决BlueSocket常见问题:超时处理、错误码解析与调试技巧
解决BlueSocket常见问题:超时处理、错误码解析与调试技巧 BlueSocket是基于Swift Package Manager的Socket框架,适用于
Elasticsearch PHP客户端版本更新深度解析
Elasticsearch PHP客户端版本更新深度解析 前言 还在为Elasticsearch PHP客户端的版本升级而头疼吗?面对从8.x到9.x的重大架构
搜索引擎后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考