Pulsar 4.2 部署、订阅可靠性与分层存储架构手册
当 Broker 健康,消费者为什么仍然收不到消息
一个常见现场是:Producer 没报错,Broker 健康检查通过,Consumer 却没有收到预期消息。真正原因可能是消息被短 topic 名写进了 public/default,消费者连接了另一个 namespace;也可能是同名 subscription 已经确认过旧消息;还可能是 Broker 返回了容器内地址,宿主机客户端完成初始连接后无法跳转到 topic owner。
Pulsar 的判断对象比“一个队列”多一层。资源名是 persistent://tenant/namespace/topic,每个 subscription 保留独立消费进度;Broker 负责协议、路由和分发,消息正文持久化在 BookKeeper,元数据存储保存 topic、ledger、bundle ownership 等协调状态。只验证 6650 端口会漏掉其中大部分故障。
排查应先锁定完整 topic 和 subscription,再确认 Broker ownership、ledger 写入、backlog 与 ack 状态。Broker 存活只是入口层证据。
从资源名字理解隔离和消费状态
Tenant 是管理与授权边界
Tenant 是多租户的第一层边界,关联允许使用的集群和管理员角色。一个业务域或团队通常拥有独立 tenant;把所有项目都放进 public,会让权限、配额、成本归属和删除责任混在一起。
Namespace 是策略落点
Namespace 位于 tenant 之下,topic 的保留、TTL、backlog quota、持久化副本、限流、隔离和授权通常在这一层配置。your_project/dev 与 your_project/prod 不只是命名差异,它们应当对应不同访问身份、容量预算和生命周期策略。
Pulsar 把 namespace 的 hash 范围拆成 bundles。Broker 获取 bundle ownership 后服务其中的 topics;Broker 扩缩容或故障时,bundle 可以卸载并由其他 Broker 接管。因此 Broker 是可横向扩展的服务层,但 bundle 热点仍可能让某个 Broker 过载。
Topic 是有分区或无分区的数据流
完整 topic 名形如:
persistent://your_project/dev/order-eventspersistent 表示消息写入 BookKeeper。分区 topic 在服务端由多个内部 partition topic 组成,Producer 可按 key 路由,Consumer 在 subscription 内并行消费。增加分区提高并行度,也增加 topic、ledger、游标和元数据数量;它不自动提高单 key 的并行度。
Subscription 才是消费视角
同一 topic 可以有多个 subscription,每个 subscription 保存独立的 cursor 和 backlog。一个 subscription 已确认消息,不影响另一个 subscription 从自己的位置消费。排障时说“topic 没消息”通常不够准确,应说“哪个 subscription 的 backlog 和 mark-delete position 是什么”。
Pulsar 支持四种常用订阅类型:
| 类型 | 消费者关系 | 顺序与故障行为 | 常见用途 |
|---|---|---|---|
| Exclusive | 仅允许一个消费者 | 最容易保持单消费者顺序;实例退出前没有并行接管 | 单实例任务、严格串行 |
| Failover | 多消费者中一个 active | active 故障后由 standby 接管;分区 topic 可在实例间分配分区 | 主备消费 |
| Shared | 多消费者共享消息 | 并行度高,不保证跨消息顺序;未确认消息可投给其他实例 | 无序工作队列 |
| Key_Shared | 相同 key 交给同一消费者 | 保持 key 内亲和与顺序,热点 key 会形成局部瓶颈 | 订单、账户等键控事件 |
选择订阅类型就是选择并行度、顺序范围和故障恢复方式。消费者数量超过分区数或可分配单元时,不一定继续提升吞吐。
Broker、BookKeeper 与元数据存储各自保存什么
Broker 是计算与协议层
Broker 接收 Producer/Consumer 连接,查找 topic owner,把写请求追加到 managed ledger,并从缓存或 BookKeeper 读取消息分发。Broker 本身不依赖本地消息盘保存业务正文,故障后可由其他 Broker 接管 bundle;但大量 bundle 同时迁移会造成连接重建、缓存冷启动和延迟波动。
Managed Ledger 把 topic 切成不可变段
Pulsar Broker 使用 managed ledger 管理 topic。一个 managed ledger 由连续的 BookKeeper ledgers 组成;当前 ledger 达到大小、时间或条目条件后封口,再创建新 ledger。不可变的已封口段让副本恢复、删除和分层存储更容易操作。
Subscription 的确认进度也会持久化。消息正文仍在 BookKeeper 中,但某个 subscription 是否需要它由 cursor 决定。所有 subscription 都确认后,消息是否继续保留由 retention 策略决定。
BookKeeper 以 E、Qw、Qa 决定副本确认
BookKeeper 集群由 bookies 组成。Broker 作为 BookKeeper client,把一个 ledger 的条目写到一组 bookies:
EnsembleSize (E):一个 ledger 选择多少个 bookies。WriteQuorum (Qw):每条 entry 写入多少个 bookies。AckQuorum (Qa):收到多少个确认后认为写入成功。
必须满足 E >= Qw >= Qa。例如 E=3, Qw=3, Qa=2 表示每条 entry 写三份,收到两份确认即可返回;这不是“永远容忍任意两节点故障”,因为恢复能力还取决于故障分布、ledger 副本位置、AutoRecovery 和剩余可写 bookies。
Bookie 使用 journal、ledger storage 和索引等本地数据。磁盘满、journal 延迟或 cookie 与数据目录不一致,都可能让 bookie 退出可写集合。Broker 健康而可用 bookies 不足时,新 ledger 无法创建,Producer 仍会失败。
元数据存储是协调状态,不是消息正文
元数据存储保存集群配置、topic 元数据、bundle ownership、BookKeeper ledger 元数据等协调信息。4.2 的常见部署仍使用 ZooKeeper;Pulsar 新架构也支持其他 metadata store,产品首页把 Oxia 列为新集群的推荐方向。技术选型必须以目标发行线和运维能力为准,不能仅替换连接串就假设迁移完成。
Standalone 把 Broker、BookKeeper 和嵌入式元数据能力放在一个 JVM 中,适合学习,不代表生产组件已经隔离。生产故障演练要分别覆盖 Broker、bookie 和元数据 quorum。
主流部署形态与生产拓扑
Standalone:一条命令跑通全部对象
Standalone 适合个人开发、SDK 验证和功能测试。优点是启动快、资源少、所有 CLI 都在一个容器中;缺点是没有组件级容灾,Broker、存储和元数据共享进程与磁盘,无法验证副本、bundle 接管或 AutoRecovery。
本地多组件 Compose:看见真实依赖关系
一个 ZooKeeper、一个 bookie、一个 Broker 的 Compose 能暴露初始化顺序、cookie、advertised listener 和 E/Qw/Qa 配置问题。它仍然没有多数派,也不能代表高可用,只适合开发者理解拓扑和自动化初始化。
单集群生产:Broker 与存储独立扩展
典型单集群包含多个 Broker、一个高可用元数据存储、多个 bookies、AutoRecovery 和统一 service URL。Broker 容量主要受连接数、协议处理、缓存、网络和 topic ownership 影响;bookie 容量主要受写入字节率、journal/ledger 磁盘、保留与积压影响。两层可独立扩展是 Pulsar 的核心优势,也意味着团队要维护两套容量模型。
官方裸机部署建议以三个 ZooKeeper 节点和三个运行 Broker/Bookie 的节点作为起步示例。生产是否共置 Broker 与 bookie,要看故障域和资源竞争:共置节省机器但 CPU、网络和页缓存互相影响;分离部署更容易独立扩容和隔离故障,成本更高。
Kubernetes / Helm:平台化交付
Helm 适合把 Broker、bookie、ZooKeeper/元数据服务、Proxy、监控和存储声明交给平台统一管理。它不会自动解决磁盘拓扑、Pod 反亲和、PDB、滚动升级或 bookie 数据恢复。StatefulSet 正常不等于 ledger 副本健康,必须同时验证 BookKeeper 和 Pulsar 层状态。
多集群与跨地域复制
多集群用于地域容灾、就近接入或数据驻留。Geo-replication 复制消息,不会自动复制所有外部业务状态,也会引入跨地域带宽、延迟、重复和故障回切问题。只有当 RPO/RTO、数据合规和演练预算确实需要时才采用,不能把它当单集群副本不足的补丁。
托管服务
托管 Pulsar 可以转移底层集群和升级责任,但 namespace 设计、subscription 生命周期、积压成本、权限角色、schema 兼容和消费幂等仍归应用团队。评审时需要核对服务端版本、协议兼容、分层存储计费、跨区流量、指标可见性和数据导出路径。
二进制 Standalone 看清运行目录
Pulsar 4.x 二进制发行版要求 Java 21。Linux 或 macOS 下载并校验 apache-pulsar-4.2.3-bin.tar.gz 后启动:
curl -LO "https://www.apache.org/dyn/closer.lua/pulsar/pulsar-4.2.3/apache-pulsar-4.2.3-bin.tar.gz?action=download"
tar xvfz apache-pulsar-4.2.3-bin.tar.gz
cd apache-pulsar-4.2.3
java -version
bin/pulsar standalone进程会创建 data/ 和 logs/。前者包含 Standalone 的 BookKeeper 与嵌入式元数据,不能在进程运行时随意删除;后者用于确认 Broker、bookie 和 metadata 初始化到了哪一步。需要后台运行时可使用:
bin/pulsar-daemon start standalone
bin/pulsar-daemon stop standaloneWindows 本机应使用 Docker 路径,不把 Linux shell 脚本改写成未经验证的批处理入口。
用容器 Standalone 从零跑通
Pulsar 下载页将 4.2.3 标为当前稳定版,5.0.0-M1 是用于提前测试的 milestone,不应作为生产基线。Standalone 容器默认使用 UID 10000、GID 0;Linux bind mount 需要相应写权限,Docker Desktop 通常由虚拟机处理映射。
docker volume create pulsar-data
docker volume create pulsar-conf
docker run -d --name pulsar-dev \
-p 127.0.0.1:6650:6650 \
-p 127.0.0.1:8080:8080 \
--mount source=pulsar-data,target=/pulsar/data \
--mount source=pulsar-conf,target=/pulsar/conf \
apachepulsar/pulsar:4.2.3 \
bin/pulsar standalone --advertised-address localhost6650 是明文 Pulsar binary protocol,8080 是 HTTP 和 Admin API。两者都只绑定回环地址,避免默认无认证的开发服务暴露到局域网。
观察启动证据:
docker logs --tail=200 pulsar-dev
docker exec pulsar-dev bin/pulsar-admin brokers healthcheck
curl -fsS http://127.0.0.1:8080/admin/v2/clusters日志应出现 messaging service ready,healthcheck 返回成功,Admin API 返回集群数据。端口已监听但日志持续出现 BookKeeper 或 metadata 初始化错误时,不要继续创建 topic;先修复 volume 权限或旧数据冲突。
用多组件 Compose 看清数据路径
Standalone 跑通后,再用一个最小多组件环境理解初始化链路。下面配置固定镜像版本,并把 E/Qw/Qa 全设为 1,因此只用于本地验证。
services:
zookeeper:
image: apachepulsar/pulsar:4.2.3
hostname: zookeeper
command: >-
bash -c "bin/apply-config-from-env.py conf/zookeeper.conf &&
bin/generate-zookeeper-config.sh conf/zookeeper.conf &&
exec bin/pulsar zookeeper"
environment:
metadataStoreUrl: zk:zookeeper:2181
volumes:
- ./data/zookeeper:/pulsar/data/zookeeper
networks: [pulsar]
healthcheck:
test: ["CMD", "bin/pulsar-zookeeper-ruok.sh"]
interval: 10s
timeout: 5s
retries: 30
pulsar-init:
image: apachepulsar/pulsar:4.2.3
command: >-
bin/pulsar initialize-cluster-metadata
--cluster cluster-a
--zookeeper zookeeper:2181
--configuration-store zookeeper:2181
--web-service-url http://broker:8080
--broker-service-url pulsar://broker:6650
depends_on:
zookeeper:
condition: service_healthy
networks: [pulsar]
bookie:
image: apachepulsar/pulsar:4.2.3
hostname: bookie
command: >-
bash -c "bin/apply-config-from-env.py conf/bookkeeper.conf &&
exec bin/pulsar bookie"
environment:
clusterName: cluster-a
metadataServiceUri: metadata-store:zk:zookeeper:2181
advertisedAddress: bookie
BOOKIE_MEM: -Xms512m -Xmx512m -XX:MaxDirectMemorySize=256m
volumes:
- ./data/bookkeeper:/pulsar/data/bookkeeper
depends_on:
zookeeper:
condition: service_healthy
pulsar-init:
condition: service_completed_successfully
networks: [pulsar]
broker:
image: apachepulsar/pulsar:4.2.3
hostname: broker
command: >-
bash -c "bin/apply-config-from-env.py conf/broker.conf &&
exec bin/pulsar broker"
environment:
metadataStoreUrl: zk:zookeeper:2181
zookeeperServers: zookeeper:2181
clusterName: cluster-a
managedLedgerDefaultEnsembleSize: 1
managedLedgerDefaultWriteQuorum: 1
managedLedgerDefaultAckQuorum: 1
advertisedAddress: broker
advertisedListeners: external:pulsar://127.0.0.1:6650
PULSAR_MEM: -Xms512m -Xmx512m -XX:MaxDirectMemorySize=256m
ports:
- "127.0.0.1:6650:6650"
- "127.0.0.1:8080:8080"
depends_on:
zookeeper:
condition: service_healthy
bookie:
condition: service_started
networks: [pulsar]
networks:
pulsar:
driver: bridge创建数据目录后启动:
mkdir -p data/zookeeper data/bookkeeper
docker compose up -d
docker compose ps
docker compose ps -a pulsar-init
docker compose logs --tail=120 zookeeper pulsar-init bookie brokerZooKeeper 必须先通过健康检查,pulsar-init 应以退出码 0 完成,Bookie 才允许启动;Broker 随后应监听 6650/8080。若初始化容器非零退出,不要靠反复重启 Broker 碰运气,先查初始化日志和 ZooKeeper 健康状态。Linux 若出现写权限错误,应按官方容器 UID 修正这两个受控目录;Windows/macOS 不要照搬 sudo chown。Bookie 日志中的 cookie mismatch 常表示数据目录和元数据属于不同集群,不能靠反复重启解决。
正向实验:创建隔离空间并证明 ack 推进
创建 tenant、namespace 和单分区 topic:
docker exec pulsar-dev bin/pulsar-admin tenants create your_project
docker exec pulsar-dev bin/pulsar-admin namespaces create your_project/dev
docker exec pulsar-dev bin/pulsar-admin topics create-partitioned-topic \
persistent://your_project/dev/order-events \
--partitions 1先启动消费者:
docker exec -it pulsar-dev bin/pulsar-client consume \
persistent://your_project/dev/order-events \
--subscription-name inventory-dev \
--subscription-type Exclusive \
--num-messages 1再生产一条带可识别内容的消息:
docker exec pulsar-dev bin/pulsar-client produce \
persistent://your_project/dev/order-events \
--messages 'event_id=evt-demo-001,status=CREATED' \
--num-produce 1消费者输出应包含同一 event_id。随后查看 topic 和 subscription:
docker exec pulsar-dev bin/pulsar-admin topics stats \
persistent://your_project/dev/order-events
docker exec pulsar-dev bin/pulsar-admin topics subscriptions \
persistent://your_project/dev/order-events有效证据包括:生产计数增加、目标 subscription 存在、msgBacklog 回到 0、消费者日志出现同一 event ID。若只看到 Producer 成功而没有 subscription,默认 retention 为关闭时消息可能在没有订阅的情况下很快变得不可依赖;需要重放能力的流应先创建 subscription 或显式设置 retention。
反向实验:制造未确认消息和地址错误
让消费者不确认,观察重新投递
应用消费者需要显式决定 ack 时机。Negative acknowledgment 适合短暂失败后的快速重投,但它的 redelivery counter 只保存在内存中,Broker 重启、bundle unload 或 Consumer 重连都可能重置;它不能可靠保证达到 maxRedeliverCount 后进入 DLQ。下面先把它用于短暂故障实验:
try (PulsarClient client = PulsarClient.builder()
.serviceUrl(System.getenv("PULSAR_SERVICE_URL"))
.build();
Consumer<byte[]> consumer = client.newConsumer()
.topic("persistent://your_project/dev/order-events")
.subscriptionName("inventory-dev")
.subscriptionType(SubscriptionType.Shared)
.negativeAckRedeliveryDelay(10, TimeUnit.SECONDS)
.subscribe()) {
while (true) {
Message<byte[]> message = consumer.receive();
try {
inventoryService.apply(message.getKey(), message.getData());
consumer.acknowledge(message);
} catch (RetryableException ex) {
consumer.negativeAcknowledge(message);
}
}
}把 inventoryService.apply 临时改为只失败一次。预期是约十秒后看到同一 message ID 或 event ID 再次出现,成功处理并 ack 后 backlog 归零。这个实验只能证明短暂重投,不能把“连续 negative ack 三次”写成持久 DLQ 保证。
需要跨重连、Broker 重启和 bundle unload 仍能累计次数时,启用 retry letter topic,并用 reconsumeLater 写入持久重试记录:
try (PulsarClient client = PulsarClient.builder()
.serviceUrl(System.getenv("PULSAR_SERVICE_URL"))
.build();
Consumer<byte[]> consumer = client.newConsumer()
.topic("persistent://your_project/dev/order-events")
.subscriptionName("inventory-dev")
.subscriptionType(SubscriptionType.Shared)
.enableRetry(true)
.deadLetterPolicy(DeadLetterPolicy.builder()
.maxRedeliverCount(3)
.retryLetterTopic("persistent://your_project/dev/order-events-retry")
.deadLetterTopic("persistent://your_project/dev/order-events-dlq")
.build())
.subscribe()) {
while (true) {
Message<byte[]> message = consumer.receive();
try {
inventoryService.apply(message.getKey(), message.getData());
consumer.acknowledge(message);
} catch (RetryableException ex) {
consumer.reconsumeLater(message, 10, TimeUnit.SECONDS);
}
}
}让指定 event ID 持续失败,第二次重试后重启 Consumer,并在另一轮隔离实验中重启 Broker 或 unload 对应 bundle。恢复后,重试次数应继续推进;超过阈值后,显式 order-events-dlq 中出现同一 event ID,主 subscription backlog 最终下降。恢复业务代码后,由受控补偿程序消费 DLQ,记录原 topic、subscription、event ID、重试次数和失败原因,再执行幂等重放。
不要把很短的 ackTimeout 当重试调度器。处理时间超过 timeout 会造成仍在执行的消息被另一消费者重新投递,形成并发重复。短暂故障可主动 negative ack;需要可持久的次数上限和最终 DLQ 时,使用 enableRetry(true) 与 reconsumeLater。
配错 advertised listener,观察两阶段连接
客户端首先连接 service URL,再根据 lookup 结果连接 topic owner。若 Broker 广播 broker:6650,宿主机客户端可能能连到入口,却无法解析或连接 broker。证据通常是初始连接成功后 lookup/redirect 超时。
从客户端实际运行环境检查:
curl -fsS http://127.0.0.1:8080/lookup/v2/topic/persistent/your_project/dev/order-events返回的 broker URL 必须从该客户端网络可达。修复 advertisedListeners 或通过 Proxy 提供稳定入口后重新 lookup;只改本机 hosts 可能掩盖集群对其他客户端仍不可达的问题。
Ack、重复、恢复和顺序的真实语义
Ack 是 subscription 的游标推进
Individual ack 确认单条消息,cumulative ack 确认当前位置及之前消息。Shared 与 Key_Shared 的并行投递不适合随意使用 cumulative ack,因为不同消费者可能仍在处理较早消息。客户端 API 和 subscription 类型共同决定允许的确认方式。
Consumer 在业务事务提交前 ack 会导致故障时消息丢失;事务成功后、ack 前退出则会重新投递。因此可靠业务通常采用至少一次消费,并在业务数据库中用 event ID 或业务版本做幂等。Pulsar 的 broker deduplication 解决 Producer 序列层面的重复,不等于业务副作用幂等。
Negative ack、ack timeout 与 redelivery
Negative ack 表示消费者明确放弃当前处理并请求稍后重投;ack timeout 表示消息在期限内没有确认。两者都可能导致消息交给另一个 Shared 消费者。客户端断开时,未确认消息也会恢复投递。
消费者恢复验收至少包含:处理前退出、业务提交后 ack 前退出、长处理超过 timeout、实例扩缩容后队列重新分配。每种情况下都要检查业务效果只有一次,backlog 最终下降,旧实例不再继续写副作用。
DLQ 与 retry letter topic
DeadLetterPolicy 为消息设置最大重试次数,并把超过阈值的消息写入 DLQ topic。DLQ 是客户端消费策略的一部分,命名和行为应由应用显式配置,不能依赖各语言客户端的默认 topic 名。Negative ack 的内存计数可能在 Broker 重启、bundle unload 和 Consumer 重连时重置;要求可靠到达 DLQ 的链路必须使用启用 retry 的 Consumer 和 reconsumeLater。
Retry letter topic 适合按计划延迟重试;negative ack 适合短暂失败。无论哪一种,都不能替代业务补偿状态机。把所有异常都重试会放大依赖故障并消耗 backlog、网络和存储。
顺序只在选定范围内成立
Exclusive 和 Failover 更容易保持单分区顺序,Shared 不保证消息顺序,Key_Shared 保证相同 key 交给同一消费者。Key_Shared 要求 Producer 提供稳定 key;热点 key 会把吞吐限制在一个消费者上。分区 topic 的跨分区全局顺序没有自然保证。
架构设计应先问业务需要“同订单有序”“同账户有序”还是“全局有序”。前两者可以用 key 和状态版本解决;全局有序会显著牺牲并行度,通常还需要下游状态机拒绝过期事件。
Retention、TTL、Backlog quota 与删除
这四个概念控制不同对象:
backlog 是某个 subscription 尚未确认的消息集合。retention 保留已被所有 subscription 确认的消息,或没有 subscription 的消息。TTL 让未确认消息超过时间后过期。
backlog quota 在未确认数据达到大小或年龄边界时执行拒写、等待或驱逐策略。
按 namespace 设置开发演示策略:
docker exec pulsar-dev bin/pulsar-admin namespaces set-retention \
your_project/dev \
--size 1G \
--time 1h
docker exec pulsar-dev bin/pulsar-admin namespaces set-backlog-quota \
your_project/dev \
--limit 512M \
--policy producer_exception
docker exec pulsar-dev bin/pulsar-admin namespaces get-retention your_project/dev
docker exec pulsar-dev bin/pulsar-admin namespaces get-backlog-quotas your_project/dev这些数值只是演示值。生产 retention 应大于 backlog quota 的相关边界,并根据峰值字节率、最长消费中断、重放窗口和磁盘预算推导。consumer_backlog_eviction 会删除尚未消费的数据,只有业务明确接受数据丢失时才能使用;producer_exception 把容量故障暴露给 Producer,应用必须有明确降级和重试上限。
清理不是即时释放磁盘。Managed ledger 先滚动并删除不再需要的 ledger,BookKeeper entry log 还要经过垃圾回收;看到 backlog 归零后,磁盘使用率可能延迟下降。判断泄漏要看多个周期的 ledger/entry log 变化,而不是立刻手工删文件。
分层存储:把封口 ledger 移到低成本介质
Pulsar 分层存储利用不可变 ledger:段封口后复制到对象存储,元数据更新指向远端位置,经过 deletion lag 后再从 BookKeeper 删除本地副本。Consumer 仍通过 Broker 读取,首次读取冷数据会增加网络延迟和对象存储请求成本。
它解决长期 backlog 和历史保留的热存储成本,不解决实时写入副本,也不是备份的同义词。对象存储凭证、bucket policy、生命周期规则、跨区流量和 offload 失败都会成为新故障面。启用前需要安装与服务端版本匹配的 offloader 包,并在所有相关 Broker 上保持一致。
分层存储的验收要同时证明:ledger 已封口并满足阈值;offload 状态成功;BookKeeper 本地空间在 deletion lag 后下降;从旧位置读取消息成功;对象存储短暂不可用时告警和重试行为符合预期。只看到对象存储里出现文件还不够。
多租户、权限与敏感数据
Pulsar 安全说明明确指出,默认没有加密、认证和授权。开发容器只能绑定本机或放在可信隔离网络;共享环境必须配置 TLS、身份认证和授权。
认证把客户端映射为 role,授权再决定 role 对 tenant、namespace 或 topic 能做什么。只启用认证而不启用授权,已认证身份仍可能访问过多资源。应用 role 只授予自己的 namespace:
bin/pulsar-admin namespaces grant-permission your_project/dev \
--actions produce,consume \
--role inventory-service
bin/pulsar-admin namespaces permissions your_project/dev更细粒度场景可分 Producer 和 Consumer role。tenant admin 可以管理 namespace,但业务应用不应使用 superuser 或 tenant admin。JWT、OAuth/OIDC、mTLS 等凭证通过密钥管理系统注入,不能进入 Git、镜像层、命令历史或 Dashboard 截图。
权限验证必须包含拒绝路径:inventory-service 可以消费 your_project/dev,访问另一个 tenant 明确失败。轮换凭证时验证长连接刷新行为,避免旧连接永久保留已撤销权限。
消息体和属性可能包含个人数据、订单号、token 与 trace baggage。Namespace 级 retention、offload bucket、备份和日志都要服从相同数据分级。删除 topic 不保证对象存储、审计日志或下游副本同步删除,数据治理要列出每个副本的 owner 和清除机制。
容量模型:Broker 与 BookKeeper 分开算
Broker 预算考虑连接数、Producer/Consumer 速率、协议线程、缓存、topic/bundle 数、网络和 JVM/direct memory。BookKeeper 预算考虑写入字节率、Qw 复制倍数、journal 与 ledger 盘吞吐、entry log GC、积压、retention、offload delay 和恢复流量。
一个可操作的存储估算起点是:
热存储 ≈ 峰值写入字节率 × 热保留时间 × Qw × 写放大安全系数
积压预算 ≈ 最大生产字节率 × 最长消费者中断时间
恢复净速率 = 恢复期消费速率 - 同期生产速率若恢复净速率不为正,扩消费者只会增加连接和重投,积压永远清不完。还要检查分区数是否足以支持并行度、业务处理是否成为瓶颈、BookKeeper 读取是否挤压在线写入。
Bookie 的磁盘阈值与只读状态必须提前告警。故障恢复会同时占用网络和磁盘,容量不能只按日常峰值填满。Rack-aware 或 zone-aware placement 要让 E/Qw/Qa 的副本真正跨故障域;三份数据落在同一宿主机或同一可用区,数字上有副本,故障上仍是单点。
恢复与故障演练
Broker 故障
Broker 退出后,bundle ownership 迁移到其他 Broker,客户端重新 lookup 并连接。通过标准是 service URL 仍可用、Producer/Consumer 自动恢复、未确认消息重新投递、业务幂等成立、延迟在预算内回落。只看到新 Broker 接管 ownership 不代表客户端已恢复。
Bookie 故障
Bookie 故障后,已有 ledger 是否可读写取决于剩余副本和 quorum;新 ledger 还需要足够 bookies 满足 E/Qw/Qa。AutoRecovery 的 auditor 和 replication worker 负责识别并补齐副本。演练应观察 under-replicated ledger、复制流量、恢复时长和故障盘重新加入时的 cookie 状态。
不要把空数据目录直接挂到旧 bookie 身份上,也不要删除 ZooKeeper/metadata store 中的 cookie 来强行启动。先确认 bookie ID、journal/ledger 目录和替换流程,否则可能同时破坏副本定位和恢复证据。
元数据 quorum 故障
元数据存储失去多数派时,已有缓存可能让部分读写短暂继续,但 ownership 变更、ledger 创建和协调操作会失败。故障演练要覆盖新建 topic、新 ledger rollover、Broker 接管和 bookie 恢复,而不是只跑一次已建立连接的 Producer。
误删与业务重放
副本和分层存储主要解决基础设施故障,不自动防止合法管理员误删 topic 或错误清 backlog。恢复方案要包含元数据备份、对象存储保护、配置即代码、审计日志和业务源数据重放。若消息是唯一事实来源,则必须单独设计不可变归档和恢复演练。
常见故障的证据顺序
消息落到 public/default
用完整 topic 列表和 Producer 配置交叉确认。短名可能解析到默认空间;修复为完整 URI,创建新 subscription 后重发测试事件。不要直接删除 public/default 中同名 topic,它可能被其他开发者使用。
同名 subscription 看不到旧消息
查看 topics stats 的 subscription cursor 和 backlog。若它已经 ack 过旧消息,创建一次性诊断 subscription 或按受控位置 reset cursor;不要把永久服务 group 随意改名,因为这会创建一份新的独立 backlog。
Bookie 启动报 cookie mismatch
先核对 metadata store URI、bookie ID、挂载目录和旧集群数据。个人环境可以在确认无共享数据后同时清理 BookKeeper 数据与对应本地元数据;共享环境必须按 bookie 恢复流程处理,不能单独删某一侧 cookie。
Producer 出现 backlog quota 异常
检查 namespace 的 quota、各 subscription backlog、最老未确认消息和消费恢复速率。若是永久消费者下线,先确认 owner 再删除 subscription;若业务需要保留,扩容存储或恢复消费。临时提高 quota 只能延后故障,还会增加恢复时间和成本。
Broker 频繁 unload bundle
观察 bundle 数、topic 热点、连接数、吞吐和 load manager 指标。热点 namespace 可以拆分 bundle 或隔离到指定 Broker,但频繁手工 unload 只会制造连接抖动。根因可能是单 key、单分区、Broker 资源不足或隔离策略错误。
项目接入与团队契约
项目配置应显式保留完整资源路径和 subscription 类型:
messaging:
pulsar:
service-url: ${PULSAR_SERVICE_URL:pulsar://127.0.0.1:6650}
admin-url: ${PULSAR_ADMIN_URL:http://127.0.0.1:8080}
topic: persistent://your_project/dev/order-events
subscription-name: inventory-dev
subscription-type: Key_Shared
ack-timeout: 60s
negative-ack-redelivery-delay: 10s
retry-enabled: true
retry-letter-topic: persistent://your_project/dev/order-events-retry
dead-letter-topic: persistent://your_project/dev/order-events-dlq
max-redeliver-count: 5
auth-plugin: ${PULSAR_AUTH_PLUGIN:}
auth-params: ${PULSAR_AUTH_PARAMS:}这些超时和次数是模板占位值,应通过业务处理时间分布、依赖恢复时间和 DLQ 处理能力确定。Consumer 还要记录 event ID、topic、subscription、message ID、redelivery count 和处理结果,消息体按数据等级脱敏。
Schema 兼容策略、partition 数、subscription 名、保留与 quota 都应由版本化配置管理。应用启动时可以验证资源存在和权限,但不要让每个实例都以管理员权限自动创建或修改 namespace 策略。
清理、回滚与资源归属
先停止 Producer 和 Consumer,确认 subscription backlog、DLQ 和 retention 数据是否仍有审计价值。删除顺序通常是临时 subscription、测试 topic、namespace、tenant;每一步都先验证 owner 和资源列表。
docker exec pulsar-dev bin/pulsar-admin topics subscriptions \
persistent://your_project/dev/order-events
docker exec pulsar-dev bin/pulsar-admin topics unsubscribe \
persistent://your_project/dev/order-events \
--subscription inventory-dev
docker exec pulsar-dev bin/pulsar-admin topics list your_project/dev
docker exec pulsar-dev bin/pulsar-admin topics delete \
persistent://your_project/dev/order-events-retry
docker exec pulsar-dev bin/pulsar-admin topics delete \
persistent://your_project/dev/order-events-dlq
docker exec pulsar-dev bin/pulsar-admin topics delete-partitioned-topic \
persistent://your_project/dev/order-events若 retry 或 DLQ topic 仍有 subscription,先列出并确认 owner,再逐一 unsubscribe;如果 topics list-partitioned-topics your_project/dev 显示它是分区 topic,则改用 delete-partitioned-topic,不能把普通 topic 的删除命令机械套用。不要用 --force 跳过归属检查。显式删除 retry/DLQ 能避免主 topic 已清理但补偿数据和费用继续残留。
个人 Standalone 环境可先保留 volume 停止:
docker stop pulsar-dev
docker start pulsar-dev确认不需要任何消息、游标和配置后再删除:
docker rm -f pulsar-dev
docker volume rm pulsar-data pulsar-conf生产回滚要分别处理 Broker 版本、bookie 数据格式、metadata schema、客户端协议和 offloader。滚动升级前先确认目标版本支持跨版本混跑,保留旧二进制与配置,验证 ledger 读写和 subscription cursor;不能通过清空 BookKeeper 或 ZooKeeper 目录“恢复干净状态”。
长期治理的判断标准
每个 tenant/namespace 应登记 owner、环境、数据等级、region/cluster、topic/partition 数、subscription、retention、TTL、backlog quota、E/Qw/Qa、offload 和成本归属。无人认领的 subscription 会持续积累数据,是最常见的隐形成本之一。
持续监控至少包括:Producer/Consumer 速率与延迟、subscription backlog 和最老消息年龄、重投与 DLQ、bundle ownership 与 unload、Broker 连接和 direct memory、bookie journal/ledger 延迟、磁盘水位、under-replicated ledgers、AutoRecovery、metadata quorum 和 offload 失败。
升级和扩容先用脱敏真实消息大小压测,再演练 Broker、bookie、metadata 和对象存储故障。通过标准是成功消息可追溯、未确认消息恢复、重复被幂等处理、backlog 按可预测斜率下降、数据副本和权限回到健康基线。
上线前逐项证明
应用始终使用完整 persistent://tenant/namespace/topic,不会误入 public/default。Subscription 类型与顺序范围、并行度和故障恢复要求一致。业务事务在 ack 之前完成,重复投递由稳定 event ID 幂等处理。
Negative ack、ack timeout、redelivery 和 DLQ 已做正反实验,并有受控重放入口。Retention、TTL 和 backlog quota 分别按重放窗口、数据丢失预算和磁盘容量配置。Broker、BookKeeper 和 metadata store 的故障模式分别演练,没有把 Standalone 当生产拓扑。
E/Qw/Qa 满足不变量,副本跨实际故障域,AutoRecovery 有容量余量。分层存储的读取、删除延迟、凭证、对象存储费用和故障告警已经验证。TLS、认证、授权和拒绝路径都已验证,业务应用不使用 superuser。
清理、升级和回滚保留 ledger、cursor、metadata 与审计证据,不靠删除数据目录恢复。
