ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Pulsar 核心术语全解析:从 Topic、订阅模式到存储与架构

Apache Pulsar 核心术语全解析:从 Topic、订阅模式到存储与架构 消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载导读本文以 Apache Pulsar 官方术语表为核心骨架系统梳理 Pulsar 中最常用、最容易混淆的 30 个核心概念覆盖消息与主题Message/Topic/Namespace/Tenant、四种订阅模式Exclusive/Shared/Failover/Key_Shared、ack/nack 消息确认机制、分布式架构组件Broker/Dispatcher/BookKeeper/配置存储以及 Pulsar Functions 与 Reader 等高级消费模型。文中所有关键概念均结合当前仓库的源码实现与配置项给出佐证读者看完不仅能读懂术语还能知道这些概念在代码中落在哪里、如何配置与使用。一、核心概念Concepts1.1 PulsarPulsar 是一个分布式发布-订阅pub-sub消息系统最初由 Yahoo 创建现由 Apache 软件基金会Apache Software Foundation管理。从当前仓库的结构可以清楚看到它的模块化组成pulsar-brokerBroker 服务端、pulsar-client-api客户端 API、pulsar-functions函数计算、pulsar-io连接器、pulsar-metadata元数据服务等各模块以 Maven 多模块工程组织在 pom.xml 中。1.2 Message消息消息Message是 Pulsar 中最基本的处理单元Producer 将消息发布到 TopicConsumer 再从 Topic 上消费消息。消息本身可以携带任意字节内容客户端 API 层通过泛型接口MessageT暴露给用户参见 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命名空间 BundleNamespace Bundle 是同一 Namespace 下的一组虚拟 Topic 集合。Bundle 被定义为一个 32 位哈希区间例如从0x00000000到0xffffffff。每个 Topic 都会被哈希映射到某个 Bundle 上而一个 Bundle 只能同时归属于一个 Broker 负责这是 Pulsar 做负载均衡的基本单元。源码佐证NamespaceBundle.java 使用 Guava 的RangeLong表示哈希区间并做了严格的边界校验区间下界必须是闭区间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 nameShared多个 Consumer 共享同一订阅名消息按 round-robin 轮流分发无the consumption order is not guaranteedFailover多个 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 中对应ProducerT接口Producer.java支持同步/异步发送、批量发送、消息去重deduplication等能力。Producer 发送消息时可指定消息 key、顺序 key、延迟投递等属性。1.11 Consumer消费者Consumer 是建立 Topic 订阅并处理 Producer 所发布消息的进程对应客户端 API 的ConsumerT接口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 Acknowledgmentack消息确认确认ack是 Consumer 发给 Pulsar Broker 的消息表示这条消息已被成功处理。ack 是 Pulsar 判断消息可以被删除或按保留策略继续保留的依据如果一条消息始终没有被确认那么它会被保留在系统中直到被处理完毕。因此 ack 直接影响消息的存储与清理。在客户端 API 中Consumer 提供多种确认方式Consumer.javaacknowledge(Message)/acknowledge(MessageId)单条确认acknowledgeCumulative(...)累积确认确认到某条消息为止的所有消息不能用于 Shared 订阅对应的异步版本acknowledgeAsync(...)以及事务内确认acknowledgeAsync(MessageId, Transaction)。1.15 Negative Acknowledgmentnack否定确认当应用程序处理某条消息失败时它可以向 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的状态。处于该状态的消息不能被删除是至少一次投递语义的基础只有收到 ackBroker 才会认为消息处理完成。1.17 Retention Policy保留策略保留策略Retention Policy是在 Namespace 上设置的大小与时间限制用于配置已被 ack 确认过的消息的保留时长/保留容量。它与 backlog未确认消息积压不同保留策略针对的是已经处理完、本可以删除的消息决定它们在满足策略前可以保留多久。配置佐证Broker 提供了默认保留配置conf/broker.conf# 默认消息保留时间分钟0 表示不保留 defaultRetentionTimeInMinutes0 # 默认保留大小MB0 表示不限制 defaultRetentionSizeInMB0 # 保留策略检查周期秒 retentionCheckIntervalInSeconds120在实际使用中可以通过pulsar-admin namespaces set-retention命令按命名空间覆盖这些默认值例如按时间保留 24 小时、按容量保留 5GB 等。1.18 Multi-Tenancy多租户多租户Multi-Tenancy是指 Pulsar 能够按租户隔离命名空间、指定配额、配置认证与授权的能力。它以 Tenant 为边界向上承载不同的业务方向下以 Namespace 进行逻辑细分是 Pulsar 在企业级部署中实现资源共享与安全隔离的核心特性。二、架构相关概念Architecture2.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、replicatorPrefixpulsar.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 BrokerBroker 是 Pulsar 集群 中的无状态组件它运行两个主要子组件HTTP Server暴露用于管理administration与主题查找topic lookup的 REST 接口Dispatcher处理所有消息传输。Pulsar 集群通常由多个 Broker 组成。Broker 本身不持久化数据——消息实际存储在 BookKeeper 中因此 Broker 可以独立扩展、故障后重启而不会丢失数据。2.9 Dispatcher分发器分发器Dispatcher是用于进出 Pulsar Broker 所有数据传输的异步 TCP 服务器所有通信都使用 Pulsar 自定义的二进制协议而非 HTTP。它承担消息分发、订阅管理、背压控制等核心运行时职责。三、存储相关概念Storage3.1 BookKeeperApache BookKeeper 是一个可扩展、低延迟的持久化日志存储服务Pulsar 用它来存储数据。Pulsar 通过 BookKeeper 获得可靠的持久化与多副本能力消息先写入 BookKeeper 的 Ledger再由 Broker 分发给消费者。3.2 BookieBookie 是单个 BookKeeper 服务器的名称从功能上看它实际上就是 Pulsar 的存储服务器。一个 BookKeeper 集群由多个 Bookie 组成消息以多副本方式默认 3 副本分布在多个 Bookie 上从而保证存储层的高可用。3.3 LedgerLedger 是 BookKeeper 中只追加append-only的数据结构用于在 Pulsar 的 Topic 上持久化存储消息。Topic 的每条消息写入 ledger 后即具备持久性只追加的特性与 BookKeeper 的分布式日志模型共同保证了顺序写的高吞吐与数据安全。与 Ledger 相关的管理、回收逻辑可从仓库的managed-ledger模块managed-ledger中进一步了解。四、函数计算FunctionsPulsar FunctionsPulsar Functions 是轻量级计算函数可以从 Pulsar Topic 消费消息、应用自定义处理逻辑并且如果需要把处理结果发布到其他 Topic。源码佐证Function.java 定义了核心接口public interface FunctionI, O { O process(I input, Context context) throws Exception; default void initialize(Context context) throws Exception {} }即用户只需实现process(input, context)即可完成输入消息 → 处理 → 输出消息的逻辑。此外还有WindowFunction.java窗口函数处理一批消息集合CollectionRecordI用于聚合、统计等窗口计算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主题可分区→ PartitionNamespace 内部按 32 位哈希切成多个 Bundle 作为负载均衡与归属分配的单元。消息生命周期Producer 发布 Message 到 Topic → BrokerDispatcher按订阅模式分发给 Consumer → Consumer 处理成功后 ack或失败时 nack / reconsumeLater 重放→ 消息在满足保留策略后被清理。消费模型Consumer 基于 Subscription四种模式消费Reader 则按任意位置扫描 Topic 消息不确认、不保留。结语Pulsar 的术语体系与它的实现是严格对应的四种订阅模式定义在 SubscriptionType.javaBundle 的哈希区间实现见 NamespaceBundle.javaack/nack 与 reconsumeLater 的能力集中在 Consumer.java保留策略的默认值则落在 conf/broker.conf。理解这些术语是阅读 Pulsar 源码、配置集群、排查消息积压与顺序问题的第一步本文可作为日常开发与运维的速查手册持续使用。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐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),仅供参考
返回列表