news 2026/10/4 10:22:04

Pulsar的Topic、Subscription和Cursors工作原理:从消息模型到消费位点管理

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Pulsar的Topic、Subscription和Cursors工作原理:从消息模型到消费位点管理

1. 从一条消息的旅程说起:Topic、Subscription 与 Cursors 到底怎么配合

如果你刚接触 Apache Pulsar,最容易懵的不是 API,而是这三个词:Topic、Subscription、Cursors。它们看起来像三个独立概念,实际上是一条消息从生产到确认的完整链路。我试过用一句话概括:Topic 是消息存放的日志,Subscription 是消费逻辑的入口,Cursors 是记录“读到哪了”的书签。理解这三者的协作,你才能搞清楚为什么 Pulsar 既能当队列用,又能当日志用。

先明确适用人群:后端开发者、消息中间件运维、以及正在做消息系统选型的人。Pulsar 的核心检索词就是 Topic 分区、Subscription 订阅类型、Cursors 消费位点管理。它和 Kafka 最大的不同在于:Pulsar 把“存储”和“消费”彻底解耦了。消息存在 Topic 里,消费进度存在 Subscription 的 Cursor 里,两者互不干扰。这意味着同一个 Topic 可以被多个订阅以不同速度、不同方式消费,而不会互相阻塞。

逻辑上,一个 Topic 就是一个追加写的日志结构,每条消息在日志里有一个偏移量(offset)。生产者把消息发到指定 Topic,Pulsar 保证消息一旦被确认(ack)就不会丢(前提是配置正确、不是整个集群挂掉)。消费者通过订阅来消费 Topic 中的消息。订阅本身不存消息数据,只存元数据和游标。游标就是那个“书签”,记录这个订阅消费到了哪个偏移量。

这里有个关键点:一个 Topic 可以挂多个订阅,每个订阅有自己独立的游标。所以订阅 A 读到 offset 100,订阅 B 可能还在 offset 20,互不影响。这就是 Pulsar 能同时支持队列语义和日志语义的底层原因——底层都是日志存储,但通过游标回放,你可以选择“消费确认后删除”(队列),也可以选择“保留并回放”(日志)。

再往下看分区。Pulsar 的分区和 Kafka 类似,但有个本质区别:Pulsar 中的分区也是 Topic。也就是说,一个分区 Topic 实际上是由多个内部 Topic 组成的,生产者可以轮询、hash 或明确指定分区来发送消息。这个设计让分区在 Pulsar 里不是特殊存在,而是 Topic 的自然延伸。

理解这三者协作,你才能回答运维中最常见的问题:为什么消息没被删除?为什么消费进度对不上?为什么共享订阅下累积确认不生效?接下来我会从环境准备开始,一步步带你创建 Topic、配置订阅、查询和重置 Cursor,最后验证消费进度。每一步都有可复制的命令和配置,你可以直接跟着做。

2. 前置准备:TaoToken 环境与 Pulsar 客户端接入配置

在动手操作之前,先把环境理顺。Pulsar 的 Topic、Subscription、Cursors 操作可以通过多种方式完成:pulsar-admin 命令行、REST API、以及各语言客户端。为了让你能快速验证,我建议用 pulsar-admin 做管理操作,用 Java/Python 客户端做消费验证。如果你还没有可用的 Pulsar 集群,可以用 Docker 起一个单机版,或者接入已有的测试集群。

这里要提一下 TaoToken。它提供统一的模型接入能力,官网是 https://taotoken.net/?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= ,API 入口是 https://taotoken.net/api 。如果你在写消费端代码时需要调用大模型做消息内容处理(比如智能路由、内容摘要),可以通过 TaoToken 的 API Key 来接入。获取 Key 的路径在控制台:https://taotoken.net/console?utm_source=taotoken_aicg_blog_end&utm_content=console&utm_campaign=rewrite ,API Keys 管理页在 https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_content=api-keys&utm_campaign=rewrite 。模型对话调试可以用 https://taotoken.net/models?utm_source=taotoken_aicg_blog_end&utm_content=models&utm_campaign=rewrite ,接入文档在 https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_content=doc&utm_campaign=rewrite 。如果你要做长期编码或 Agent 类任务,可以看 Coding Plan:https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_content=coding-plan&utm_campaign=rewrite 。Claude Code 相关接入参考 https://taotoken.net/claude-code?utm_source=taotoken_aicg_blog_end&utm_content=claude-code&utm_campaign=rewrite 。

回到 Pulsar。先确认你的 pulsar-admin 能连上集群。假设你的集群服务地址是pulsar://localhost:6650,Web 服务地址是http://localhost:8080。你可以用环境变量或配置文件指定。下面是一个典型的客户端配置片段,以 Java 客户端为例,放在pulsar-client.properties或代码里的PulsarClient.builder()中:

# pulsar-client.properties serviceUrl=pulsar://localhost:6650 operationTimeoutMs=30000 connectionTimeoutMs=10000

如果你用 Python 客户端,配置类似:

import pulsar client = pulsar.Client( 'pulsar://localhost:6650', operation_timeout_seconds=30, connection_timeout_ms=10000 )

对于 pulsar-admin,通常通过--url指定 Web 服务地址:

pulsar-admin --url http://localhost:8080 topics list public/default

这里public/default是租户/命名空间。Pulsar 的 Topic 全名格式是persistent://tenant/namespace/topic。默认租户是public,默认命名空间是default。你可以先列出已有 Topic 确认连通性。

如果你用的是 TaoToken 的 API 来做消费端增强,比如在消费到消息后调用模型做分类,那么你需要在客户端代码里配置 Base URL 和 Key。以 OpenAI 兼容方式为例:

from openai import OpenAI client = OpenAI( base_url="https://taotoken.net/api", api_key="你的TaoToken API Key" ) response = client.chat.completions.create( model="gpt-4o-mini", messages=[{"role": "user", "content": "对这条消息做意图分类:..."}] )

注意,TaoToken 的 API 入口是https://taotoken.net/api,不要加 UTM 参数到 API 地址上。Key 从控制台获取。这样你就能在 Pulsar 消费逻辑里嵌入模型调用,实现智能处理。

环境准备好后,我们进入核心操作:创建 Topic、配置订阅、管理 Cursor。每一步我都会给出命令和预期结果。

3. 可复制配置:Topic 创建、订阅类型设置与 Cursor 位点管理

这一节是全文的核心操作区。我会按“创建 Topic → 创建订阅 → 查询 Cursor → 重置 Cursor”的顺序,给出可直接复制的命令和配置片段。你可以在测试集群上跟着做。

3.1 创建分区 Topic

先创建一个带分区的 Topic。假设我们要创建一个 3 分区的持久化 Topic,名为my-topic,在public/default命名空间下:

pulsar-admin topics create-partitioned-topic \ persistent://public/default/my-topic \ --partitions 3

执行后,Pulsar 会创建 3 个内部分区 Topic:my-topic-partition-0、my-topic-partition-1、my-topic-partition-2。你可以用以下命令确认:

pulsar-admin topics list-partitioned-topics public/default pulsar-admin topics partitions persistent://public/default/my-topic

预期输出会列出 3 个分区。注意,分区 Topic 本身不存储消息,消息实际存在各分区里。生产者发送时如果不指定分区,默认轮询。

3.2 创建订阅并指定类型

Pulsar 支持四种订阅类型:Exclusive(独享)、Shared(共享)、Failover(灾备)、Key_Shared(键共享)。创建订阅时指定类型:

pulsar-admin topics create-subscription \ persistent://public/default/my-topic \ --subscription my-sub \ --subscription-type Shared

这里创建了一个名为my-sub的共享订阅。共享订阅下可以有多个消费者同时消费,消息在消费者间竞争分发。如果你要严格顺序,用 Exclusive;如果要主备切换,用 Failover;如果要按 key 保证顺序且多消费者,用 Key_Shared。

创建后可以查看订阅列表:

pulsar-admin topics subscriptions persistent://public/default/my-topic

预期输出包含my-sub及其类型。

3.3 查询 Cursor 位点

Cursor 是订阅的消费位点。查询某个订阅的 Cursor 位置:

pulsar-admin topics peek-messages \ persistent://public/default/my-topic \ --subscription my-sub \ --count 1

或者用更直接的方式查看订阅统计:

pulsar-admin topics stats persistent://public/default/my-topic

在 stats 输出中,找到subscriptions字段,里面有msgBacklog(积压消息数)、msgRateOut、msgThroughputOut等。msgBacklog为 0 表示该订阅已消费完当前所有消息。你还可以用:

pulsar-admin topics stats-internal persistent://public/default/my-topic

这个命令会显示每个分区的cursor信息,包括markDeletePosition(已确认删除位置)和readPosition(当前读取位置)。markDeletePosition就是 Cursor 的核心位点。

3.4 重置 Cursor 位点

重置 Cursor 是运维常用操作,比如消费出错需要回滚重放。命令如下:

pulsar-admin topics reset-cursor \ persistent://public/default/my-topic \ --subscription my-sub \ --message-id 10:5:-1

--message-id格式是ledgerId:entryId:partitionIndex。你也可以用时间戳重置:

pulsar-admin topics reset-cursor \ persistent://public/default/my-topic \ --subscription my-sub \ --time 1h

--time 1h表示重置到 1 小时前。重置后,该订阅会从指定位置重新消费。注意,重置 Cursor 不会删除消息,只是移动书签。

3.5 配置文件片段(settings/TOML/JSON)

如果你用 Pulsar 的配置文件方式管理订阅,可以在pulsar-admin的配置或客户端 settings 中写入。以下是一个 JSON 格式的订阅配置示例,用于客户端初始化:

{ "topic": "persistent://public/default/my-topic", "subscriptionName": "my-sub", "subscriptionType": "Shared", "receiverQueueSize": 1000, "ackTimeoutMillis": 30000, "negativeAckRedeliveryDelayMillis": 60000 }

如果你用 TOML 管理(比如某些运维脚本),可以写成:

[topic] name = "persistent://public/default/my-topic" partitions = 3 [subscription] name = "my-sub" type = "Shared" ack_timeout_ms = 30000

这些配置片段可以直接放进你的项目 settings 文件或客户端初始化代码中。注意路径和原文一致,不要随意改字段名。

3.6 消费端代码示例(Java)

下面是一个 Java 消费者示例,展示如何用 Shared 订阅消费并确认消息:

import org.apache.pulsar.client.api.*; public class MyConsumer { public static void main(String[] args) throws Exception { PulsarClient client = PulsarClient.builder() .serviceUrl("pulsar://localhost:6650") .build(); Consumer<byte[]> consumer = client.newConsumer() .topic("persistent://public/default/my-topic") .subscriptionName("my-sub") .subscriptionType(SubscriptionType.Shared) .ackTimeout(30, java.util.concurrent.TimeUnit.SECONDS) .subscribe(); while (true) { Message<byte[]> msg = consumer.receive(); try { System.out.println("收到消息: " + new String(msg.getValue())); // 处理消息 consumer.acknowledge(msg); } catch (Exception e) { consumer.negativeAcknowledge(msg); } } } }

这段代码创建了一个 Shared 订阅的消费者,收到消息后确认。如果处理失败,用negativeAcknowledge让消息重新投递。注意,Shared 模式下累积确认不适用,但可以用批量确认减少 RPC 调用。

3.7 累积确认与批量确认

Pulsar 支持单条确认和累积确认。累积确认吞吐量更高,但失败时会重复处理。在 Exclusive 或 Failover 订阅下可以用:

consumer.acknowledgeCumulative(msg);

Shared 模式下不能用累积确认,但可以用批量确认:

consumer.acknowledgeAsync(msg);

或者用consumer.acknowledgeCumulativeAsync在支持的订阅类型下异步确认。

以上配置和命令覆盖了 Topic 创建、订阅类型设置、Cursor 查询与重置。接下来我们验证请求是否成功,并检查消费进度。

4. 验证请求与成功结果:消费进度检查与位点确认

配置完成后,必须验证消息链路是否按预期工作。这一节我会给出生产消息、消费消息、检查 Cursor 位点的完整验证步骤,并说明每个步骤的成功标志。

4.1 生产测试消息

先用 pulsar-client 或命令行生产几条消息。如果你有 pulsar-client 工具:

pulsar-client produce \ persistent://public/default/my-topic \ --messages "msg-1" "msg-2" "msg-3" \ --num-produce 1

预期输出类似:

2025-01-01 10:00:00 INFO [ProducerImpl] Producer created 2025-01-01 10:00:00 INFO [ProducerImpl] Published 3 messages

如果看到Published 3 messages,说明生产成功。你也可以用 Java 生产者:

Producer<byte[]> producer = client.newProducer() .topic("persistent://public/default/my-topic") .create(); for (int i = 1; i <= 3; i++) { producer.send(("msg-" + i).getBytes()); } producer.close();

4.2 消费消息并确认

用上一节的消费者代码启动消费。启动后,控制台应输出:

收到消息: msg-1 收到消息: msg-2 收到消息: msg-3

每收到一条,调用acknowledge确认。确认后,Cursor 的markDeletePosition会前移。

4.3 检查 Cursor 位点

消费确认后,再次查询订阅统计:

pulsar-admin topics stats persistent://public/default/my-topic

在输出中找到subscriptions下的my-sub,关注msgBacklog。如果为 0,说明所有消息已确认。同时看msgRateOut和msgThroughputOut,确认有消费流量。

再用stats-internal查看具体位点:

pulsar-admin topics stats-internal persistent://public/default/my-topic

输出中每个分区会有:

"cursor": { "my-sub": { "markDeletePosition": "10:5:-1", "readPosition": "10:6:-1" } }

markDeletePosition是已确认删除的位置,readPosition是当前读取位置。如果markDeletePosition接近readPosition,说明消费进度正常。

4.4 验证消息保留与删除

如果你没有设置保留策略,当所有订阅的 Cursor 都消费到某个偏移量后,该偏移量之前的消息会被自动删除。你可以用:

pulsar-admin topics stats-internal persistent://public/default/my-topic

查看ledgers列表。如果某个 ledger 的所有消息都被确认,它会被标记为可删除。你可以用:

pulsar-admin topics expire-messages \ persistent://public/default/my-topic \ --subscription my-sub

手动触发过期。但注意,这需要所有订阅都确认。

4.5 验证重置 Cursor

重置 Cursor 后,再次消费应该从指定位置开始。例如:

pulsar-admin topics reset-cursor \ persistent://public/default/my-topic \ --subscription my-sub \ --message-id 10:0:-1

然后重启消费者,应该重新收到msg-1开始的消息。如果收到重复消息,说明重置成功。

4.6 成功结果汇总

一个完整的成功验证链路应该是:

  1. 生产 3 条消息,Published 3 messages。
  2. 消费者收到 3 条并确认。
  3. msgBacklog为 0。
  4. markDeletePosition前移。
  5. 重置 Cursor 后能重新消费。

如果你在消费端集成了 TaoToken 的模型调用,比如对消息做分类,那么验证时还要确认模型返回正常。你可以用模型对话页面 https://taotoken.net/models?utm_source=taotoken_aicg_blog_end&utm_content=models&utm_campaign=rewrite 先调试好 prompt,再放进消费代码。

验证通过后,我们来看常见错误和排查方法。

5. 本篇常见错排查:401、local proxy failed、reading choices、OAuth 报错对照

在实际操作中,你可能会遇到各种报错。这一节我整理了几类高频错误,包括 Pulsar 本身的报错和接入 TaoToken 时的报错,给出原因和解决步骤。

5.1 Pulsar 订阅相关报错

报错:SubscriptionNotFoundException

org.apache.pulsar.client.api.PulsarClientException$SubscriptionNotFoundException: Subscription my-sub not found

原因:订阅名写错,或者订阅还没创建。解决:先用pulsar-admin topics subscriptions确认订阅存在。如果不存在,用create-subscription创建。

报错:ConsumerBusyException

org.apache.pulsar.client.api.PulsarClientException$ConsumerBusyException: Exclusive consumer is already connected

原因:Exclusive 订阅下已经有消费者连接,第二个消费者被拒绝。解决:改用 Shared 或 Failover 订阅,或者断开已有消费者。

报错:TopicNotFoundException

Topic persistent://public/default/my-topic not found

原因:Topic 未创建,或者命名空间不对。解决:用pulsar-admin topics create-partitioned-topic创建,确认租户/命名空间。

5.2 Cursor 重置报错

报错:InvalidMessageIdException

Invalid message id: 10:5:-1

原因:message-id 格式不对,或者指定的 ledger/entry 不存在。解决:用stats-internal查看有效的 ledgerId 和 entryId,确保格式是ledgerId:entryId:partitionIndex。

报错:CursorResetException

Failed to reset cursor: subscription is not active

原因:订阅没有活跃消费者,或者集群状态异常。解决:先启动一个消费者,再重置;或者检查 broker 日志。

5.3 TaoToken 接入报错

报错:401 Unauthorized

{"error": {"message": "Invalid API key", "type": "invalid_request_error"}}

原因:API Key 错误或未设置。解决:检查api_key是否从 https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_content=api-keys&utm_campaign=rewrite 正确获取。注意 Base URL 是https://taotoken.net/api,不要多加路径。

报错:local proxy failed

Error: local proxy failed: connection refused

原因:本地代理配置问题,或者网络不通。解决:检查你的 HTTP 客户端是否配置了代理,确保能访问https://taotoken.net/api。如果你在容器里跑,确认 DNS 和出网正常。

报错:reading choices

KeyError: 'choices'

原因:模型返回结构不符合预期,通常是请求格式不对或模型名错误。解决:检查model参数是否在 TaoToken 支持的模型列表里,参考 https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_content=doc&utm_campaign=rewrite 。确保请求体是标准的 chat completions 格式。

报错:OAuth 相关

OAuth token expired or invalid

原因:如果你用 OAuth 方式接入,token 过期。解决:重新获取 token,或者改用 API Key 方式。TaoToken 的 API Key 方式更简单,直接在控制台生成即可。

5.4 消费进度不推进

现象:msgBacklog一直不降

原因可能有三:消费者没有确认消息;确认超时;或者订阅类型不匹配。解决:检查代码里是否调用了acknowledge;检查ackTimeout设置;确认 Shared 订阅下没有用累积确认。

现象:消息重复消费

原因:确认失败或超时,消息被重新投递。解决:确保处理逻辑幂等;调整ackTimeout;用negativeAcknowledge显式重投。

5.5 配置三件套检查

如果你在 Pulsar 消费端集成了 TaoToken,或者用 Cline MCP、Codex auth.json 等方式接入,务必检查三件套:Base URL、Key、Model ID。以 Codex 的auth.json为例:

{ "base_url": "https://taotoken.net/api", "api_key": "你的Key", "model": "gpt-4o-mini" }

Cline MCP 配置类似,在 settings 里填好 Base URL、Key、Model ID。CC Switch 也是同样三件套。缺一不可,否则会报 401 或 model not found。

排查完这些,你的 Pulsar 消息链路应该能稳定运行了。最后说一下后续怎么继续深入。

6. 继续深入:从 Cursor 位点管理到生产级消费链路

走到这里,你已经掌握了 Topic 创建、Subscription 类型选择、Cursor 查询与重置、消费进度验证,以及常见报错排查。但生产环境还有几个点值得继续打磨。

第一,保留策略与 Cursor 的配合。默认情况下,所有订阅确认后消息会被删除。但如果你设置了保留策略(按时间或大小),已确认的消息会保留到阈值再删除。这在需要回放历史消息的场景很有用。你可以用:

pulsar-admin namespaces set-retention public/default \ --size 10G \ --time 7d

这样即使所有订阅都确认了,消息还会保留 7 天或 10G。Cursor 重置后可以回放。

第二,多订阅协同。一个 Topic 可以挂多个订阅,每个订阅独立消费。比如一个订阅做实时处理,一个订阅做离线分析。它们的 Cursor 互不影响。你可以用stats查看每个订阅的积压情况,分别扩容消费者。

第三,Key_Shared 订阅的顺序保证。如果你需要按消息 key 保证顺序,同时又要多消费者并行,用 Key_Shared。配置时指定keySharedPolicy:

consumer.newConsumer() .topic("persistent://public/default/my-topic") .subscriptionName("my-sub") .subscriptionType(SubscriptionType.Key_Shared) .keySharedPolicy(KeySharedPolicy.autoSplitHashRange()) .subscribe();

第四,消费端集成模型处理。如果你在消费消息后需要调用大模型做内容理解,可以用 TaoToken 的 API。先通过模型对话页面调试 prompt,再放进消费代码。长期跑编码或 Agent 任务的话,Coding Plan 更合适:https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_content=coding-plan&utm_campaign=rewrite 。接入文档在 https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_content=doc&utm_campaign=rewrite ,API Key 在 https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_content=api-keys&utm_campaign=rewrite 。

第五,监控 Cursor 滞后。生产环境要监控msgBacklog和markDeletePosition与readPosition的差值。如果积压持续增长,说明消费能力不足,需要扩容消费者或优化处理逻辑。你可以用 Prometheus + Grafana 采集 Pulsar 的 stats 指标。

最后,记住一个原则:Cursor 是订阅的,不是 Topic 的。重置 Cursor 只影响该订阅,不影响其他订阅。删除订阅会删除其 Cursor,但不会删除消息。理解这一点,你就能灵活设计消费链路。

如果你在实操中遇到其他报错,可以先查接入文档,或者在模型对话页面用 TaoToken 调试你的处理逻辑。消息中间件的运维没有银弹,多动手验证,多观察 Cursor 位点变化,慢慢就有感觉了。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/4 10:19:58

MRAM+8位MCU实战:MR25H40CDF与PIC18F45K50的高可靠工业存储设计

1. 这个组合能做什么&#xff1a;MR25H40CDF 与 PIC18F45K50 的应用背景前一阵在调一块工业采集板&#xff0c;主控是 Microchip 的 PIC18F45K50&#xff0c;数据存储从原来的 SPI EEPROM 换成了 Everspin 的 MR25H40CDF。项目需求很典型&#xff1a;现场设备要记录参数修改、事…

作者头像 李华
网站建设 2026/10/4 10:19:10

OpenClaw为什么叫“龙虾”?附本地部署与API Key配置详解

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华