ARTICLE DETAIL

资讯详情

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

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

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。生产者把消息发到指定 TopicPulsar 保证消息一旦被确认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_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content API 入口是 https://taotoken.net/api 。如果你在写消费端代码时需要调用大模型做消息内容处理比如智能路由、内容摘要可以通过 TaoToken 的 API Key 来接入。获取 Key 的路径在控制台https://taotoken.net/console?utm_sourcetaotoken_aicg_blog_endutm_contentconsoleutm_campaignrewrite API Keys 管理页在 https://taotoken.net/api-keys?utm_sourcetaotoken_aicg_blog_endutm_contentapi-keysutm_campaignrewrite 。模型对话调试可以用 https://taotoken.net/models?utm_sourcetaotoken_aicg_blog_endutm_contentmodelsutm_campaignrewrite 接入文档在 https://taotoken.net/doc?utm_sourcetaotoken_aicg_blog_endutm_contentdocutm_campaignrewrite 。如果你要做长期编码或 Agent 类任务可以看 Coding Planhttps://taotoken.net/coding-plan?utm_sourcetaotoken_aicg_blog_endutm_contentcoding-planutm_campaignrewrite 。Claude Code 相关接入参考 https://taotoken.net/claude-code?utm_sourcetaotoken_aicg_blog_endutm_contentclaude-codeutm_campaignrewrite 。回到 Pulsar。先确认你的 pulsar-admin 能连上集群。假设你的集群服务地址是pulsar://localhost:6650Web 服务地址是http://localhost:8080。你可以用环境变量或配置文件指定。下面是一个典型的客户端配置片段以 Java 客户端为例放在pulsar-client.properties或代码里的PulsarClient.builder()中# pulsar-client.properties serviceUrlpulsar://localhost:6650 operationTimeoutMs30000 connectionTimeoutMs10000如果你用 Python 客户端配置类似import pulsar client pulsar.Client( pulsar://localhost:6650, operation_timeout_seconds30, connection_timeout_ms10000 )对于 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_urlhttps://taotoken.net/api, api_key你的TaoToken API Key ) response client.chat.completions.create( modelgpt-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 个内部分区 Topicmy-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(); Consumerbyte[] 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) { Messagebyte[] 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 生产者Producerbyte[] 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 成功结果汇总一个完整的成功验证链路应该是生产 3 条消息Published 3 messages。消费者收到 3 条并确认。msgBacklog为 0。markDeletePosition前移。重置 Cursor 后能重新消费。如果你在消费端集成了 TaoToken 的模型调用比如对消息做分类那么验证时还要确认模型返回正常。你可以用模型对话页面 https://taotoken.net/models?utm_sourcetaotoken_aicg_blog_endutm_contentmodelsutm_campaignrewrite 先调试好 prompt再放进消费代码。验证通过后我们来看常见错误和排查方法。5. 本篇常见错排查401、local proxy failed、reading choices、OAuth 报错对照在实际操作中你可能会遇到各种报错。这一节我整理了几类高频错误包括 Pulsar 本身的报错和接入 TaoToken 时的报错给出原因和解决步骤。5.1 Pulsar 订阅相关报错报错SubscriptionNotFoundExceptionorg.apache.pulsar.client.api.PulsarClientException$SubscriptionNotFoundException: Subscription my-sub not found原因订阅名写错或者订阅还没创建。解决先用pulsar-admin topics subscriptions确认订阅存在。如果不存在用create-subscription创建。报错ConsumerBusyExceptionorg.apache.pulsar.client.api.PulsarClientException$ConsumerBusyException: Exclusive consumer is already connected原因Exclusive 订阅下已经有消费者连接第二个消费者被拒绝。解决改用 Shared 或 Failover 订阅或者断开已有消费者。报错TopicNotFoundExceptionTopic persistent://public/default/my-topic not found原因Topic 未创建或者命名空间不对。解决用pulsar-admin topics create-partitioned-topic创建确认租户/命名空间。5.2 Cursor 重置报错报错InvalidMessageIdExceptionInvalid message id: 10:5:-1原因message-id 格式不对或者指定的 ledger/entry 不存在。解决用stats-internal查看有效的 ledgerId 和 entryId确保格式是ledgerId:entryId:partitionIndex。报错CursorResetExceptionFailed 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_sourcetaotoken_aicg_blog_endutm_contentapi-keysutm_campaignrewrite 正确获取。注意 Base URL 是https://taotoken.net/api不要多加路径。报错local proxy failedError: local proxy failed: connection refused原因本地代理配置问题或者网络不通。解决检查你的 HTTP 客户端是否配置了代理确保能访问https://taotoken.net/api。如果你在容器里跑确认 DNS 和出网正常。报错reading choicesKeyError: choices原因模型返回结构不符合预期通常是请求格式不对或模型名错误。解决检查model参数是否在 TaoToken 支持的模型列表里参考 https://taotoken.net/doc?utm_sourcetaotoken_aicg_blog_endutm_contentdocutm_campaignrewrite 。确保请求体是标准的 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。配置时指定keySharedPolicyconsumer.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_sourcetaotoken_aicg_blog_endutm_contentcoding-planutm_campaignrewrite 。接入文档在 https://taotoken.net/doc?utm_sourcetaotoken_aicg_blog_endutm_contentdocutm_campaignrewrite API Key 在 https://taotoken.net/api-keys?utm_sourcetaotoken_aicg_blog_endutm_contentapi-keysutm_campaignrewrite 。第五监控 Cursor 滞后。生产环境要监控msgBacklog和markDeletePosition与readPosition的差值。如果积压持续增长说明消费能力不足需要扩容消费者或优化处理逻辑。你可以用 Prometheus Grafana 采集 Pulsar 的 stats 指标。最后记住一个原则Cursor 是订阅的不是 Topic 的。重置 Cursor 只影响该订阅不影响其他订阅。删除订阅会删除其 Cursor但不会删除消息。理解这一点你就能灵活设计消费链路。如果你在实操中遇到其他报错可以先查接入文档或者在模型对话页面用 TaoToken 调试你的处理逻辑。消息中间件的运维没有银弹多动手验证多观察 Cursor 位点变化慢慢就有感觉了。
返回列表