Apache Pulsar 与云原生消息中间件

共 20 题
📑 题目列表 20 题
#
★★★

1. Pulsar 的分层架构,Broker 无状态与 BookKeeper 存储分离的设计动机

Pulsar 的分层架构(Broker 无状态与 BookKeeper 存储分离)的设计动机是什么?

  • Broker 无状态
  • BookKeeper 存储层
  • 设计动机(弹性/扩展/高可用)

Pulsar 采用"存储计算分离":Broker 无状态只负责协议处理与读写调度,数据持久化到 BookKeeper(独立存储集群)。设计动机:一是弹性扩展——Broker 无状态可水平增减,扩容不需数据搬迁;二是存储与计算解耦——BookKeeper 负责多副本持久化,Broker 负载与存储分离,写放大与读放大可分别优化;三是高可用——Broker 故障时其他 Broker 接管,BookKeeper 副本保证数据不丢;四是多租户与长保留——存储层独立支持分层存储(卸载到对象存储)。对比 Kafka 的存储计算耦合,Pulsar 度独立扩展更灵活,适合云原生与弹性伸缩。

存储计算分离让 Broker 无状态、可独立扩容,BookKeeper 提供多副本持久化与分层存储,是 Pulsar 弹性、高可用、多租户的基础。

#
★★★

2. Pulsar 与 Kafka 的核心差异,存储计算分离与日志追加模型对比

Pulsar 与 Kafka 的核心差异(存储计算分离与日志追加模型)是什么?

  • 存储计算分离
  • 日志追加模型
  • 架构差异

Pulsar 与 Kafka 核心差异:Pulsar 存储计算分离(Broker 无状态 + BookKeeper 存储),Kafka 存储计算耦合(Broker 本地磁盘日志)。日志追加模型上,两者都是"追加日志",但 Pulsar 用分段(Segment/ledger)追加到 BookKeeper,Kafka 用分区日志追加到本地磁盘。差异影响:Pulsar 扩容无须数据搬迁、支持分层存储与长保留、多租户隔离好;Kafka 分区固定副本在 broker、扩缩容需 rebalance(数据搬迁)、长保留受磁盘限制。Pulsar 更适合云原生弹性、多租户、长保留;Kafka 生态成熟、与大数据/流处理集成好。选型按对弹性、多租户、保留的需求。

核心差异是"存储计算分离 vs 耦合",影响扩展性、保留能力与多租户。日志追加模型两者都有,但底层存储不同。

#
★★★

3. Pulsar 的多租户模型,tenant/namespace/topic 三级模型与配额治理

Pulsar 的多租户模型(tenant/namespace/topic 三级模型)与配额治理是什么?

  • tenant/namespace/topic 三级
  • 多租户隔离
  • 配额治理

Pulsar 多租户用"tenant/namespace/topic"三级模型:tenant(租户,如一个部门/团队)-> namespace(命名空间,如一个应用/环境,可设策略)-> topic(主题)。租户与命名空间隔离:不同 tenant 有独立认证、权限、配额,namespace 可配置保留策略、最大消息大小、backlog 配额等。配额治理:通过 namespace 级策略限制(如吞吐、存储、消息数)、权限(生产/消费)、持久化策略,实现资源隔离与治理。多租户让多个业务共享集群但互不干扰,且按 namespace 管理策略,是 Pulsar 云原生/多租户的核心能力。

三级模型实现"租户隔离 + namespace 策略 + topic 消息",配额治理通过 namespace 级限制与权限控制实现资源隔离。

#
★★★

4. Pulsar 的消费模型,共享/独占/故障转移/key-shared 四种订阅类型

Pulsar 的消费模型(共享/独占/故障转移/key-shared 四种订阅类型)是什么?

  • 独占订阅
  • 共享订阅
  • 故障转移与 key-shared

Pulsar 四种订阅类型:Exclusive(独占)——一个 topic 只被一个消费者订阅,其他消费者被拒,保证严格顺序;Shared(共享)——多个消费者分摊消息(round-robin),无顺序保证,高吞吐;Failover(故障转移)——一个主消费者,其他为备,主消费者故障时备接管,保留顺序;Key_Shared(键共享)——按 key 哈希把消息路由到固定消费者,同一 key 的消息由同一消费者处理(保序 + 并行)。选型:要严格顺序用 Exclusive,要并行高吞吐用 Shared,要高可用且保序用 Failover,要"按 key 保序 + 并行"用 Key_Shared。

四种订阅对应"顺序 vs 并行"的权衡:Exclusive 严格顺序、Shared 高并行、Failover 高可用保序、Key_Shared 按 key 保序并行。

#
★★★

5. Pulsar 的消息确认与重试,ack 超时、负确认(nack)与死信主题

Pulsar 的消息确认与重试(ack 超时、负确认 nack 与死信主题)如何工作?

  • ack 确认
  • ack 超时与 nack
  • 死信主题

Pulsar 消费者处理成功后 acknowledge 确认消息;若未确认,消息在 ackTimeout(ack 超时)后重新投递(至少一次)。negativeAcknowledge(nack)显式告诉 broker 消息处理失败,可重投(可设 nackRedeliveryDelayMs)。重试次数用尽或消息持续失败可投递到死信主题(deadLetterTopic),供人工处理/告警。机制配合:ack 成功移除、ack 超时/nack 重投、死信兜底,实现 at-least-once + 失败隔离。acknowledgeCumulative 可批量确认(顺序消费)。Java 端用 Consumer.acknowledge/negativeAcknowledge 与配置。

ack 超时与 nack 触发重投,死信主题隔离最后失败,保证 at-least-once 与失败可治理。

#
★★★

6. Pulsar Reader API,从指定 messageId 读取、与 Consumer 的差异,以及数据回放与复制管道的用途

Pulsar Reader API 如何从指定 messageId 读取,与 Consumer 的差异,以及数据回放与复制管道的用途是什么?

  • Reader API 概念
  • 与 Consumer 的差异
  • 回放与复制管道

Pulsar Reader API 从指定 messageId(或 earliest/latest)读取,可自由定位并逐条读取,不依赖消费订阅/ack 机制(Reader 无需确认,适用于读取型场景)。与 Consumer 的差异:Consumer 面向"消费组 + 订阅 + ack + 重投"(消息处理后确认),Reader 面向"指定位置读取 + 无 ack 语义",适合数据回放(从历史 offset 重读)、复制管道(把 topic 数据复制到其他系统/存储)、迁移与审计。用途:数据回放(重放历史消息)、复制管道(topic -> 下游/存储)、读取特定位置数据。Reader 是"拉取式只读"接口,与 Consumer 的"消费式"互补。

Reader 是"指定位置拉取 + 无 ack"的只读接口,适合回放与复制,Consumer 是"消费组 + ack"的消费接口,适合作业处理。

#
★★

7. Pulsar 的延迟消息与定时消息实现与精度限制

Pulsar 的延迟消息与定时消息实现与精度限制是什么?

  • 延迟消息原理
  • 定时消息
  • 精度限制

Pulsar 延迟消息通过 MessageBuilder.deliverAt 指定投递时间,broker 在消息到达时若不达投递时间,放入"延迟消息队列"(delay deliver tracker),由 broker 定时检查(默认每 1 秒)到期后投递。精度限制:延迟投递由 broker 的定时扫描触发,精度约秒级(扫描间隔 1s),非毫秒级;延迟消息在 broker 内存中跟踪(可持久化),大量延迟消息影响内存。Java 端用 newMessage().deliverAt(time) 发送延迟消息。适合秒级延迟(如定时任务、限时业务),不适合毫秒级精确。对比 Kafka 需自扩展,Pulsar 内置延迟消息。

Pulsar 延迟消息靠 broker 定时扫描(约 1s 精度),精度秒级,内置支持,适合定时/延迟场景。

#
★★

8. Pulsar Functions 与连接器(connector)生态的工程应用

Pulsar Functions 与连接器(connector)生态的工程应用是什么?

  • Pulsar Functions 轻量计算
  • 连接器(source/sink)
  • 工程应用

Pulsar Functions 是轻量级流处理函数(类似 Kafka Streams/Lambda),在 topic 上执行 process 函数处理消息(过滤、转换、聚合),可部署为独立函数,无需单独流处理框架。连接器(connector)分 Source(把外部数据源写入 Pulsar,如 Kafka、数据库、文件)与 Sink(把 Pulsar 消息写到外部系统,如数据库、对象存储、日志)。工程应用:用 Functions 做消息清洗/转换/报警,用 Source/Sink 与其他系统集成(数据管道),两者都基于 Pulsar 生态,通过函数式与连接器实现数据集成与流处理,减少自建。Java 实现 Function 接口或 Source/Sink 类。

Functions 提供轻量计算,连接器提供与外部系统集成,两者构成 Pulsar 的流处理与数据管道生态。

#
★★

9. Pulsar 事务(消息事务)与跨主题原子性保证

Pulsar 事务(消息事务)与跨主题原子性保证是什么?

  • Pulsar 事务
  • 跨主题原子
  • 事务协调器

Pulsar 事务(pulsar-client 事务 API)让消息写入多个 topic、以及消费确认(ack)在同一事务内原子提交,通过事务协调器(transaction coordinator)管理,read_committed 消费者只读已提交事务消息。跨主题原子性:一个事务可原子写入多个 topic 的消息,要么全部可见要么全部不可见,配合消费 ack 实现"生产 + 消费"的端到端原子(exactly-once 语义的基础)。Pulsar 事务依赖 transactionCoordinatorEnabled 与事务补偿。Java 端用 TransactionBuilder/client.newTransaction() 启动事务。适合需要跨主题一致性的场景(如事件溯源 + 聚合)。

Pulsar 事务通过事务协调器实现跨主题原子写入与消费确认,是 exactly-once 与跨主题一致性的基础。

#
★★

10. Pulsar 的分层存储(tiered storage),卸载到 S3/GCS 的机制与成本

Pulsar 的分层存储(tiered storage)卸载到 S3/GCS 的机制与成本是什么?

  • 分层存储机制
  • 卸载到对象存储
  • 成本与收益

Pulsar 分层存储(tiered storage):topic 数据先写入 BookKeeper(热数据),达到阈值/时间后,后台把旧的 Segment 卸载(offload)到对象存储(S3/GCS/Azure Blob),消费者读旧数据时从对象存储读取。机制:Segment 是卸载单位,BookKeeper 保留元数据,对象存储存数据。成本:热数据占 BookKeeper(贵、SSD),冷数据卸载到对象存储(便宜、大容量),降低存储成本并支持海量数据长期保留;但读取冷数据延迟更高(对象存储网络)。收益:长保留 + 低成本;权衡:冷数据读取延迟。适合需要长期保留历史数据(审计、事件溯源)的场景。

分层存储把热数据放 BookKeeper、冷数据卸载到对象存储,用"分层 + 卸载"平衡成本与长期保留,冷读延迟略高。

#
★★

11. Pulsar Java 客户端的生产者/消费者 API 与批量/压缩配置

Pulsar Java 客户端的生产者/消费者 API 与批量/压缩配置是什么?

  • ProducerBuilder/ConsumerBuilder
  • 批量发送与压缩
  • 配置项

Pulsar Java 客户端用 PulsarClient 创建 ProducerBuilder(生产者)与 ConsumerBuilder(消费者)。生产者配置:topicproducerNamecompressionType(压缩:LZ4/ZSTD/ZLIB)、batchingEnabled(批量)、batchingMaxMessages/batchingMaxPublishDelay(批量阈值)、sendTimeout;消费者配置:subscriptionType(订阅类型)、ackTimeoutreceiverQueueSize(接收队列)、subscriptionName。批量发送:启 batchingEnabled 后多条消息合并发送,减吞吐;压缩:compressionType 减少网络带宽。Java 端用 producer.newMessage()/producer.send() 发送,consumer.receive()/acknowledge() 消费。

Pulsar 客户端 Builder 配置生产者/消费者,批量(batching)与压缩(compression)提升吞吐与带宽效率。

#
★★

12. Pulsar 的 exactly-once 语义,消息去重与幂等生产者

Pulsar 的 exactly-once 语义如何通过消息去重与幂等生产者实现?

  • 幂等生产者
  • 消息去重
  • exactly-once

Pulsar 实现 exactly-once 依靠:幂等生产者(producerName 唯一 + 序号去重)+ 事务(跨生产/消费原子)+ 消息去重(broker 按消息序号去重)。幂等生产者:同一 producer 重试发送时通过序号去重,避免重复写入;配合事务可把"生产 + 消费"原子化。exactly-once 是"生产端幂等 + 传输不重 + 消费端去重"的组合:生产端幂等去重、broker 去重、消费端幂等处理。Pulsar 的幂等与去重依赖事务与消息序号,保证每条消息恰好处理一次。Java 端启用事务与幂等配置实现。

exactly-once 是"幂等生产 + 事务 + 去重"的组合,Pulsar 用生产者序号去重与事务机制实现,消费端需幂等配合。

#
★★

13. Pulsar 与 Spring Boot 集成(spring-pulsar starter)的注解式消费

Pulsar 与 Spring Boot 集成(spring-pulsar starter)的注解式消费如何实现?

  • spring-pulsar starter
  • @PulsarListener 注解
  • 配置与模板

spring-pulsar starter 提供 Spring Boot 对 Pulsar 的集成:PulsarTemplate(发送)、PulsarListenerContainerFactory(消费容器)、@PulsarListener(注解式消费)与 @PulsarReader(读取)。注解式消费:@PulsarListener(subscriptionName = "...", topics = "...") 标注方法,Spring 自动创建消费者容器并消费,方法参数为消息(可配 @Payload@Header)。配置:spring.pulsar.client.*(客户端)、spring.pulsar.producer.*spring.pulsar.consumer.*(如订阅类型、ack 超时)。可扩展 PulsarListenerContainerFactory 定制。对比 Kafka 的 @KafkaListener,Pulsar 用 @PulsarListener,声明式消费 + 客户端模板。

spring-pulsar 提供 PulsarTemplate 发送 + @PulsarListener 注解消费,配合 properties 配置,与 Spring 生态集成。

#
★★

14. Pulsar 集群的高可用,BookKeeper 写入链路、ack quorum 与读修复

Pulsar 集群的高可用(BookKeeper 写入链路、ack quorum 与读修复)如何实现?

  • BookKeeper 写入链路
  • ack quorum
  • 读修复

Pulsar 高可用依赖 BookKeeper 的写入链路:写入时 broker 把 Segment 复制到多个 Bookie(ack quorum),达到多数派(quorum)确认后才算成功,保证数据不丢;ensemble size(写副本数)与 ack quorum(多数派确认)配置可靠性。读修复:读请求向多个 Bookie 读取,若某副本数据损坏/不一致,用多数派正确数据修复(read repair),保证数据一致性。故障处理:Bookie 故障时,Pulsar 自动替换/恢复,Segment 按 quorum 复制保证可用。高可用 = 多副本 + 多数派确认 + 读修复 + Bookie 故障恢复。

Pulsar 高可用靠 BookKeeper 多副本(ack quorum)+ 多数派确认 + 读修复,broker 无状态故障迁移,保障数据不丢与可用。

#
★★

15. Pulsar 的消息积压治理,消费 backlog 监控、扩容与跳过策略

Pulsar 的消息积压治理(消费 backlog 监控、扩容与跳过策略)如何做?

  • backlog 监控
  • 扩容消费
  • 跳过策略

Pulsar 积压治理:一是监控 backlog(topicbacklogbacklogQuota),用 pulsar-admin 或监控系统查看消息积压量;二是扩容消费——增加消费者(Shared/Key_Shared 订阅可并行,或增加 subscription 的消费者),提升消费并行度;三是跳过策略——对无法处理/积压严重的历史消息用 seek 跳过或丢弃(pulsar-adminskip/clean),必要时清零 backlog;四是优化消费(批量、异步、下游优化)。配合 backlogQuota 设置积压上限(超过后丢弃/拒绝)。治理目标:监控 backlog 定位、扩容提吞吐、跳过/丢弃处理过期积压。

积压治理是"监控 backlog + 扩容消费 + 跳过/丢弃 + 优化消费"的组合,backlogQuota 设上限防无限积压。

#
★★

16. Pulsar 的 schema registry,schema 版本演进与兼容性策略(FORWARD/BACKWARD/FULL),如何避免不兼容变更

Pulsar 的 schema registry 的版本演进与兼容性策略(FORWARD/BACKWARD/FULL)如何避免不兼容变更?

  • schema 版本演进
  • 兼容性策略
  • 避免不兼容

Pulsar schema registry 管理消息 schema(Avro/JSON/Protobuf),每个 topic 绑定 schema 并支持版本演进。兼容性策略:BACKWARD(新 schema 能读旧数据,即消费者升级后能读旧消息)、FORWARD(新数据能被旧消费者读,即先升级生产)、FULL(同时向后和向前兼容)。Pulsar 在 schema 变更时校验兼容性,不兼容变更(如删除字段、改类型)会被拒绝。避免不兼容:字段只能新增(带默认值)、不能删除/改类型、按兼容性策略演进,用 schema 版本管理。这样生产/消费者可陆续升级而不破坏。Java 端用 Schema.AVRO 等与 schema 注册。

schema 演进需遵守兼容性策略(BACKWARD/FORWARD/FULL),只增字段不删不改,Pulsar 校验阻止不兼容变更。

#
★★

17. Pulsar 的跨地域复制(geo-replication),基于 topic 的复制方向、集群拓扑与异步复制的一致性边界

Pulsar 的跨地域复制(geo-replication)基于 topic 的复制方向、集群拓扑与异步复制的一致性边界是什么?

  • 跨地域复制基于 topic
  • 复制方向与拓扑
  • 异步复制一致性边界

Pulsar 跨地域复制(geo-replication)基于 topic 级别配置:一个 topic 可配置复制到多个集群,事件发布到任一集群后,异步复制到其他集群(发布-订阅式复制)。复制方向:可单向(一主多从)或双向(多集群互备),按 topic 配置复制集群列表。集群拓扑:多集群间通过跨地域连接(inter-cluster replication)异步复制。一致性边界:复制是异步的,存在延迟(跨地域延迟),因此不同集群看到的数据可能短暂不一致(最终一致);复制顺序按分区保证,跨集群无强一致。适合多活/容灾/就近读,但不适合强一致需求。Pulsar 也可用"pulsar 集群联邦"配置。

跨地域复制是 topic 级异步复制,方向与拓扑可配置,一致性是最终一致(异步延迟),适合多活容灾非强一致。

#

18. Pulsar 的负载均衡与 Topic 归属(ownership)动态迁移

Pulsar 的负载均衡与 Topic 归属(ownership)动态迁移如何工作?

  • Topic 归属
  • 负载均衡
  • 动态迁移

Pulsar 中每个 topic 的读写由某个 broker 拥有(ownership),broker 无状态,topic 归属可动态迁移。负载均衡:Pulsar 的负载均衡器(LoadShedding/LoadManager)监控 broker 负载,把过载 broker 上的 topic 迁移到低负载 broker(topic 归属迁移),实现负载均衡。动态迁移:topic 归属变化时,broker 通过元数据服务(ZK)协调接管,客户端自动感知(重连新 broker)。迁移不影响数据(数据在 BookKeeper),只迁移计算责任,因此迁移快、无数据搬迁。适合弹性伸缩与负载均衡。副作用:迁移期间 topic 可能短暂不可用(几十 ms)。

topic 归属 + 负载均衡器实现动态迁移,因为数据在 BookKeeper,迁移只搬计算责任,快且无数据搬迁。

#

19. Pulsar 与 Kafka/RocketMQ 在云原生部署(K8s/Operator)下的运维对比

Pulsar 与 Kafka/RocketMQ 在云原生部署(K8s/Operator)下的运维对比是什么?

  • 云原生部署
  • Operator 管理
  • 运维对比

云原生(K8s)部署下:Pulsar 有 streamnative/pulsar-operator,支持 K8s 原生部署(StatefulSet + 无状态 broker),存储计算分离使 broker 可弹性伸缩、BookKeeper 用 PVC 持久化,Operator 管理生命周期/扩缩容/升级,云原生适配好。Kafka 用 Strimzi/Confluent Operator 管理 Kafka 集群(broker + KRaft),RocketMQ 用 rocketmq-operator。对比:Pulsar 存储计算分离使 broker 无状态、弹性扩缩容更顺;Kafka broker 有状态(本地数据),扩缩容需 rebalance;RocketMQ 有 NameServer/Broker。运维复杂度上,Pulsar 依赖 BookKeeper(组件多、运维复杂),Kafka 依赖 KRaft(相对简单),RocketMQ 依赖 NameServer。云原生部署三者都支持 Operator,但 Pulsar 的弹性与多租户更适配 K8s。

云原生下三者都支持 Operator,Pulsar 存储计算分离弹性更好但组件(BookKeeper)运维复杂,Kafka 相对简单,RocketMQ 中等。

#

20. Pulsar 的认证与授权,JWT/OAuth 认证流程、namespace 级生产消费权限与 TLS 传输加密

Pulsar 的认证与授权(JWT/OAuth 认证流程、namespace 级生产消费权限与 TLS 传输加密)如何实现?

  • JWT/OAuth 认证
  • namespace 级权限
  • TLS 加密

Pulsar 认证支持 JWT、OAuth2、TLS 认证等。JWT 认证:客户端携带签名的 JWT token,broker 验证签名后建立身份;OAuth2:用 OAuth2 授权码/客户端凭证获取 token 访问。授权:Pulsar 按 namespace/tenant 级别授予权限(produce/consume/functions),通过 pulsar-admin 或 API 给角色授权,实现 namespace 级生产/消费权限控制。TLS:客户端与 broker 之间用 TLS 传输加密(tlsEnabled、证书配置),保护数据在传输中的机密性。组合:认证(验证身份)+ 授权(控制权限)+ TLS(传输加密),实现安全访问。Java 端配置 AuthenticationToken/OAuth 与 TLS 证书。

安全 = 认证(JWT/OAuth 验证身份)+ 授权(namespace 级权限)+ 传输加密(TLS),三者配合实现安全消息访问。