- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
导读
本文以 Apache Pulsar 官方术语表为核心骨架,系统梳理 Pulsar 中最常用、最容易混淆的 30+ 个核心概念,覆盖消息与主题(Message/Topic/Namespace/Tenant)、四种订阅模式(Exclusive/Shared/Failover/Key_Shared)、ack/nack 消息确认机制、分布式架构组件(Broker/Dispatcher/BookKeeper/配置存储)以及 Pulsar Functions 与 Reader 等高级消费模型。文中所有关键概念均结合当前仓库的源码实现与配置项给出佐证,读者看完不仅能"读懂"术语,还能知道这些概念在代码中落在哪里、如何配置与使用。
一、核心概念(Concepts)
1.1 Pulsar
Pulsar 是一个分布式发布-订阅(pub-sub)消息系统,最初由 Yahoo 创建,现由 Apache 软件基金会(Apache Software Foundation)管理。从当前仓库的结构可以清楚看到它的模块化组成:pulsar-broker(Broker 服务端)、pulsar-client-api(客户端 API)、pulsar-functions(函数计算)、pulsar-io(连接器)、pulsar-metadata(元数据服务)等,各模块以 Maven 多模块工程组织在 pom.xml 中。
1.2 Message(消息)
消息(Message)是 Pulsar 中最基本的处理单元,Producer 将消息发布到 Topic,Consumer 再从 Topic 上消费消息。消息本身可以携带任意字节内容,客户端 API 层通过泛型接口Message<T>暴露给用户(参见 Message.java)。
1.3 Topic(主题)
Topic 是一个具名的通道,用于把 Producer 发布的消息传递给负责处理这些消息的 Consumer。在 Pulsar 中,Topic 的完整命名通常遵循persistent://tenant/namespace/topic或non-persistent://tenant/namespace/topic的格式,其中persistent/non-persistent前缀决定了消息是否落盘持久化。
1.4 Partitioned Topic(分区主题)
分区主题(Partitioned Topic)是由多个 Pulsar Broker 共同提供服务的主题,从而获得更高的吞吐能力。分区本质上把一条逻辑主题拆分成多个物理分片,每个分区可以独立路由到不同的 Broker,从源码结构看,分区的路由、管理与订阅逻辑集中在pulsar-client-api的PartitionedProducerImpl/PartitionedTopicImpl等类中。要注意的是,术语"分区"与下文 Namespace Bundle 是两回事:分区是消息流层面的水平拆分,Bundle 是命名空间在负载均衡层面的拆分。
1.5 Namespace(命名空间)
命名空间(Namespace)是相关 Topic 的分组机制,是 Pulsar 多租户体系中的管理单元之一。命名空间的命名通常形如tenant/namespace,可以理解为"租户下的逻辑分组"。Pulsar 允许对命名空间级别的 Topic 统一配置保留策略、限流、权限等。
1.6 Namespace Bundle(命名空间 Bundle)
Namespace Bundle 是同一 Namespace 下的一组虚拟 Topic 集合。Bundle 被定义为一个 32 位哈希区间,例如从0x00000000到0xffffffff。每个 Topic 都会被哈希映射到某个 Bundle 上,而一个 Bundle 只能同时归属于一个 Broker 负责,这是 Pulsar 做负载均衡的基本单元。
源码佐证:NamespaceBundle.java 使用 Guava 的Range<Long>表示哈希区间,并做了严格的边界校验:
- 区间下界必须是闭区间(
BoundType.CLOSED); - 区间上界除非等于全上界
0xffffffff,否则必须是开区间(BoundType.OPEN),即不允许两个 Bundle 的哈希范围重叠; - Bundle 的字符串形式由
String.format("0x%08x_0x%08x", ...)生成,例如0x00000000_0xffffffff。
当某个 Bundle 负载过高时,Pulsar 会将其分裂(bundle split)成两个更小的区间,由负载均衡器重新分配到不同的 Broker。
1.7 Tenant(租户)
租户(Tenant)是用于分配容量、实施认证(authentication)与授权(authorization)方案的管理单元。Pulsar 的多租户能力以租户为顶层:管理员按租户分配配额,并为每个租户配置独立的认证授权策略。租户-命名空间-Topic 的层级关系构成了 Pulsar 资源隔离的骨架。
1.8 Subscription(订阅)与四种订阅模式
订阅(Subscription)是消费组在 Topic 上建立的一种"租约"(lease),由一组 Consumer 建立。Pulsar 共有四种订阅模式:exclusive、shared、failover 与 key_shared,它们定义在客户端 API 的枚举中,见 SubscriptionType.java:
| 订阅模式 | 特点 | 顺序保证 | 源码注释要点 |
|---|---|---|---|
Exclusive | 同一个订阅名下只允许 1 个 Consumer | 有 | "There can be only 1 consumer on the same topic with the same subscription name" |
Shared | 多个 Consumer 共享同一订阅名,消息按 round-robin 轮流分发 | 无 | "the consumption order is not guaranteed" |
Failover | 多个 Consumer 使用同一订阅名,但同一时刻只有 1 个活跃 Consumer 接收消息;该 Consumer 断开后由其他 Consumer 接管 | 有 | 分区主题下按"每个分区最多 1 个活跃 Consumer"的方式拆分分区分配,顺序按分区粒度保证 |
Key_Shared | 多个 Consumer 共享同一订阅,相同 key 的消息只分发给同一个 Consumer | 按 key 保证 | 可通过ordering_key覆盖消息 key 以影响排序 |
在 Shared 模式下,多个 Consumer 可以并行消费以提升吞吐,但全局顺序不被保证;Failover 与 Key_Shared 则在不同粒度上兼顾了顺序与并发。实际使用中,可通过ConsumerBuilder#subscriptionType(...)指定模式。
1.9 Pub-Sub(发布-订阅)
发布-订阅是一种消息传递模式:Producer 进程把消息发布到 Topic 上,然后由 Consumer 进程消费(处理)这些消息。发布者与订阅者之间通过 Topic 解耦,彼此无需知晓对方的存在。
1.10 Producer(生产者)
Producer 是向 Pulsar Topic 发布消息的进程。在客户端 API 中对应Producer<T>接口(Producer.java),支持同步/异步发送、批量发送、消息去重(deduplication)等能力。Producer 发送消息时可指定消息 key、顺序 key、延迟投递等属性。
1.11 Consumer(消费者)
Consumer 是建立 Topic 订阅并处理 Producer 所发布消息的进程,对应客户端 API 的Consumer<T>接口(Consumer.java)。Consumer 通过receive()拉取消息、通过acknowledge()向 Broker 确认处理完成,也可通过negativeAcknowledge()通知 Broker 重放消息(详见下文 ack/nack)。
1.12 Reader(读取器)
Reader 是另一类消息处理器,与 Consumer 非常相似,但有两个关键差异:
- 可以指定从 Topic 的哪个位置开始处理消息(而 Consumer 总是从最新的未确认消息开始);
- Reader 不保留数据,也不确认(ack)消息——它像"游标扫描"一样按位置读取消息,适合消息回放、状态重建、审计等场景。
源码佐证:Reader.java 的接口注释明确指出 "A Reader can be used to scan through all the messages currently available in a topic",并提供readNext()与带超时的readNext(int timeout, TimeUnit unit)等方法,可通过MessageId.earliest / latest或任意指定消息 ID 定位读取起点。
1.13 Cursor(游标)
游标(Cursor)是某个 Consumer 的订阅位置(subscription position)。它记录该消费者在订阅中已经读到哪条消息,Broker 依据游标决定下次分发的起点。
1.14 Acknowledgment(ack,消息确认)
确认(ack)是 Consumer 发给 Pulsar Broker 的消息,表示"这条消息已被成功处理"。ack 是 Pulsar 判断消息可以被删除(或按保留策略继续保留)的依据:如果一条消息始终没有被确认,那么它会被保留在系统中,直到被处理完毕。因此 ack 直接影响消息的存储与清理。
在客户端 API 中,Consumer 提供多种确认方式(Consumer.java):
acknowledge(Message)/acknowledge(MessageId):单条确认;acknowledgeCumulative(...):累积确认(确认到某条消息为止的所有消息,不能用于 Shared 订阅);- 对应的异步版本
acknowledgeAsync(...),以及事务内确认acknowledgeAsync(MessageId, Transaction)。
1.15 Negative Acknowledgment(nack,否定确认)
当应用程序处理某条消息失败时,它可以向 Pulsar 发送"否定确认"(negative ack),通知系统稍后重放这条消息。默认情况下,被 nack 的消息会在1 分钟延迟后被重放。
重要提示:在有序订阅类型(Exclusive、Failover、Key_Shared)上使用 nack,可能导致失败消息以乱序的方式重新到达消费者——因为被 nack 的消息会跳过其后的有序消息重新投递。
源码佐证:Consumer.java 中negativeAcknowledge(Message)的注释明确说明其重投延迟可通过ConsumerBuilder#negativeAckRedeliveryDelay(long, TimeUnit)配置;此外还提供了reconsumeLater(msg, delayTime, unit)方法,允许对单条消息指定自定义延迟后重新投递。二者都让失败消息的"重放时机"变得可控制。
1.16 Unacknowledged(未确认)
未确认(unacknowledged)指消息已被投递给 Consumer 进行处理、但尚未被 Consumer 确认(ack)的状态。处于该状态的消息不能被删除,是"至少一次投递"语义的基础:只有收到 ack,Broker 才会认为消息处理完成。
1.17 Retention Policy(保留策略)
保留策略(Retention Policy)是在 Namespace 上设置的大小与时间限制,用于配置已被 ack 确认过的消息的保留时长/保留容量。它与 backlog(未确认消息积压)不同:保留策略针对的是"已经处理完、本可以删除"的消息,决定它们在满足策略前可以保留多久。
配置佐证:Broker 提供了默认保留配置(conf/broker.conf):
# 默认消息保留时间(分钟),0 表示不保留 defaultRetentionTimeInMinutes=0 # 默认保留大小(MB),0 表示不限制 defaultRetentionSizeInMB=0 # 保留策略检查周期(秒) retentionCheckIntervalInSeconds=120在实际使用中,可以通过pulsar-admin namespaces set-retention命令按命名空间覆盖这些默认值,例如按时间保留 24 小时、按容量保留 5GB 等。
1.18 Multi-Tenancy(多租户)
多租户(Multi-Tenancy)是指 Pulsar 能够按租户隔离命名空间、指定配额、配置认证与授权的能力。它以 Tenant 为边界,向上承载不同的业务方,向下以 Namespace 进行逻辑细分,是 Pulsar 在企业级部署中实现资源共享与安全隔离的核心特性。
二、架构相关概念(Architecture)
2.1 Standalone(单机模式)
Standalone 是一种轻量级 Pulsar Broker:所有组件都运行在单个 JVM 进程中。Standalone 集群可以在单台机器上运行,主要用于开发调试。仓库中的 conf/standalone.conf 即为单机模式的默认配置,覆盖了 Broker、BookKeeper、ZooKeeper 等组件的本地化参数。
2.2 Cluster(集群)
集群(Cluster)是一组 Pulsar Broker 与 BookKeeper 服务器(即 Bookie)的集合。不同地理区域的集群之间可以通过 Geo-Replication 相互复制消息。集群是 Pulsar 部署的基本单元,一个实例(Instance)通常由多个集群组成。
2.3 Instance(实例)
实例(Instance)是一组协同工作、作为一个整体单元的 Pulsar 集群。多集群实例常见于需要跨地域容灾或全局统一的逻辑部署场景。
2.4 Geo-Replication(地理复制)
地理复制(Geo-Replication)指消息在多个 Pulsar 集群 之间的复制,这些集群可能位于不同的数据中心或地理区域。它让消息能够跨地域同步,是 Pulsar 支撑全球部署的关键能力。相关配置见 conf/broker.conf 中replicationConnectionsPerBroker、replicationProducerQueueSize、replicatorPrefix=pulsar.repl等复制相关参数。
2.5 Configuration Store(配置存储)
配置存储(Configuration Store)是 Pulsar 用于配置类任务的ZooKeeper quorum(注意:术语表原文将其描述为"previously known as configuration store",即早期版本中它曾被称作"configuration store")。多集群 Pulsar 安装只需要一个跨所有集群共享的配置存储,用于保存全局配置元数据。在现代版本中,Pulsar 也支持使用 etcd 等作为元数据存储(见pulsar-metadata模块)。
2.6 Topic Lookup(主题查找)
主题查找(Topic Lookup)是 Pulsar Broker 提供的服务:让连接的客户端自动确定某个 Topic 由哪个 Pulsar 集群负责,从而把该 Topic 的消息流量路由到正确位置。从源码看,Broker 的二进制协议处理链路中实现了查找命令处理,例如 ServerCnx.java 中的handleLookup(CommandLookupTopic)即为查找请求的服务端入口之一。
2.7 Service Discovery(服务发现)
服务发现(Service Discovery)是 Pulsar 提供的一种机制:连接的客户端只需使用一个 URL 就能与集群中的所有 Broker 交互。客户端把请求发送到统一入口,由服务发现层将其引导到正确的 Broker(通常经由 Topic Lookup 完成),从而屏蔽了集群内部多 Broker 的复杂性。
2.8 Broker
Broker 是 Pulsar 集群 中的无状态组件,它运行两个主要子组件:
- HTTP Server:暴露用于管理(administration)与主题查找(topic lookup)的 REST 接口;
- Dispatcher:处理所有消息传输。
Pulsar 集群通常由多个 Broker 组成。Broker 本身不持久化数据——消息实际存储在 BookKeeper 中,因此 Broker 可以独立扩展、故障后重启而不会丢失数据。
2.9 Dispatcher(分发器)
分发器(Dispatcher)是用于进出 Pulsar Broker 所有数据传输的异步 TCP 服务器,所有通信都使用 Pulsar 自定义的二进制协议(而非 HTTP)。它承担消息分发、订阅管理、背压控制等核心运行时职责。
三、存储相关概念(Storage)
3.1 BookKeeper
Apache BookKeeper 是一个可扩展、低延迟的持久化日志存储服务,Pulsar 用它来存储数据。Pulsar 通过 BookKeeper 获得可靠的持久化与多副本能力,消息先写入 BookKeeper 的 Ledger,再由 Broker 分发给消费者。
3.2 Bookie
Bookie 是单个 BookKeeper 服务器的名称,从功能上看,它实际上就是 Pulsar 的存储服务器。一个 BookKeeper 集群由多个 Bookie 组成,消息以多副本方式(默认 3 副本)分布在多个 Bookie 上,从而保证存储层的高可用。
3.3 Ledger
Ledger 是 BookKeeper 中只追加(append-only)的数据结构,用于在 Pulsar 的 Topic 上持久化存储消息。Topic 的每条消息写入 ledger 后即具备持久性;只追加的特性与 BookKeeper 的分布式日志模型共同保证了顺序写的高吞吐与数据安全。与 Ledger 相关的管理、回收逻辑可从仓库的managed-ledger模块(managed-ledger)中进一步了解。
四、函数计算(Functions)
Pulsar Functions
Pulsar Functions 是轻量级计算函数:可以从 Pulsar Topic 消费消息、应用自定义处理逻辑,并且(如果需要)把处理结果发布到其他 Topic。
源码佐证:Function.java 定义了核心接口:
public interface Function<I, O> { O process(I input, Context context) throws Exception; default void initialize(Context context) throws Exception {} }即用户只需实现process(input, context)即可完成"输入消息 → 处理 → 输出消息"的逻辑。此外还有:
- WindowFunction.java:窗口函数,处理一批消息集合
Collection<Record<I>>,用于聚合、统计等窗口计算; - Context.java:向执行中的函数提供上下文信息(如当前 Topic、日志、状态存储等)。
Pulsar Functions 与普通 Consumer 相比的优势在于:它把"消费 → 处理 → 产出"封装成无运维负担的轻量计算单元,天然与 Topic 模型集成。仓库的pulsar-functions模块还提供了 Java/Python/Go 等多种语言的运行时与大量示例(见 pulsar-functions)。
五、概念关系速查
把上述概念串起来,Pulsar 的整体数据流与部署层级可以这样理解:
- 部署层级:Instance(实例)→ Cluster(集群)→ Broker + Bookie;多集群之间通过 Geo-Replication 复制,跨集群共享一个 Configuration Store。
- 资源层级:Tenant(租户)→ Namespace(命名空间)→ Topic(主题,可分区)→ Partition;Namespace 内部按 32 位哈希切成多个 Bundle 作为负载均衡与归属分配的单元。
- 消息生命周期:Producer 发布 Message 到 Topic → Broker(Dispatcher)按订阅模式分发给 Consumer → Consumer 处理成功后 ack(或失败时 nack / reconsumeLater 重放)→ 消息在满足保留策略后被清理。
- 消费模型:Consumer 基于 Subscription(四种模式)消费;Reader 则按任意位置扫描 Topic 消息,不确认、不保留。
结语
Pulsar 的术语体系与它的实现是严格对应的:四种订阅模式定义在 SubscriptionType.java,Bundle 的哈希区间实现见 NamespaceBundle.java,ack/nack 与 reconsumeLater 的能力集中在 Consumer.java,保留策略的默认值则落在 conf/broker.conf。理解这些术语,是阅读 Pulsar 源码、配置集群、排查消息积压与顺序问题的第一步;本文可作为日常开发与运维的速查手册持续使用。
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar 核心术语全解:从消息、主题到多租户与存储架构
Apache Pulsar 核心术语全解:从消息、主题到多租户与存储架构 本文是 Apache Pulsar 官方术语表( site2/website next
消息队列后端流处理Apache Pulsar 术语大全:从 Message、Topic 到 Broker、BookKeeper 的核心概念体系解析
Apache Pulsar 术语大全:从 Message、Topic 到 Broker、BookKeeper 的核心概念体系解析 Apache Pulsar 是
消息队列后端流处理Apache Pulsar 术语表:核心概念、架构组件与存储原理解析
Apache Pulsar 术语表:核心概念、架构组件与存储原理解析 本指南以 Apache Pulsar 官方术语文档为主体,系统梳理从消息、主题、命名空间、
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考