Apache Pulsar 消息队列实践:通过 Shared 订阅与 Receiver Queue 将 Topic 用作消息队列

发布时间:2026/9/24 16:43:31
Apache Pulsar 消息队列实践:通过 Shared 订阅与 Receiver Queue 将 Topic 用作消息队列 Apache Pulsar 消息队列实践通过 Shared 订阅与 Receiver Queue 将 Topic 用作消息队列【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsar导读本文围绕 Apache Pulsar 官方文档中「Using Pulsar as a message queue」这一实战主题展开讲解如何在不引入额外中间件的情况下把 Pulsar 的 Topic 直接当作消息队列使用。你将掌握共享订阅Shared subscription与接收队列receiver queue两个核心配置的组合技巧并得到 Java、Python、C、Go 四种语言客户端的可直接运行示例以及对应的源码级原理验证。这套方案适合「每条工作对象都必须被处理、即使某个组件变慢或宕机也不能丢失数据」的强一致任务分发场景。在许多大规模数据架构中消息队列是不可或缺的组件。当系统中每一个工作对象都必须被处理——即使某个组件处理缓慢甚至完全失败——就需要一个消息队列来兜底确保未处理的数据按正确顺序被保留直到完成所需的处理动作。Apache Pulsar 天然适合承担这一角色原因有二Pulsar 从设计之初就以持久化消息存储为核心消息落盘于 BookKeeper可抵御 broker 重启与订阅方故障切换Pulsar 支持在 Topic 的消费者之间自动进行消息负载均衡如果你愿意也可以自定义负载均衡策略。一份部署两种用法你可以在同一个 Pulsar 集群上同时充当实时消息总线message bus与消息队列message queue或者只使用其中一种。可以划出部分 Topic 专供实时场景、另一部分 Topic 专供消息队列场景如果需要更严格的隔离也可以用不同的命名空间namespace分别承载两类用途。为什么 Pulsar 适合做消息队列Pulsar 的消息传递模型基于发布-订阅模式。理解该模型是正确配置消息队列的前提生产者producer将消息发布到 Topic消费者consumer通过**订阅subscription**接收消息创建订阅后Pulsar 会保留所有消息即使消费者断开连接也不会丢失只有消费者确认ack消息处理成功后消息才会被标记为可删除订阅分为 exclusive独占、shared共享、failover故障转移、key_shared按键共享四种类型。要获得消息队列的「多消费者分摊处理」语义必须使用shared类型。从源码结构可以进一步确认实现「消息队列」效果的关键是让多个消费者使用相同的订阅名即共享同一条订阅shared、failover、key_shared 均属此列而若每个消费者使用各自唯一的订阅名则退化为传统的扇出发布-订阅pub-sub模式。客户端配置改造让 Topic 变成消息队列要把 Pulsar Topic 用作消息队列需要在消费者侧做两件事1. 建立共享订阅Shared subscription所有消费者必须建立 shared 订阅并使用完全相同的订阅名。否则订阅便不是共享的各消费者无法组成处理集群processing ensemble消息也就无法在多个消费者之间分摊。2. 压低接收队列receiver queue大小如果希望在各消费者之间精确控制消息分发可以把消费者的接收队列大小设置得很低必要时甚至可以设为 0。每个 Pulsar 消费者都有一个接收队列决定消费者一次尝试预取多少条消息。例如接收队列为1000默认值时消费者连接后一次会尝试从 Topic 积压backlog中拉取最多 1000 条消息进行本地缓冲接收队列设为0或极小值时基本可以确保每个消费者同一时刻只处理一件事从而对消息在消费者之间的分派做到最精细的控制。在客户端配置数据类中可以看到两个默认值的源码定义配置项默认值源码位置subscriptionTypeSubscriptionType.Exclusive独占ConsumerConfigurationData.javareceiverQueueSize1000ConsumerConfigurationData.java这意味着默认情况下 Pulsar 消费者是独占订阅 1000 条预取队列距离消息队列的语义还有两步配置要走——改订阅类型为 Shared、按需调低接收队列。限制接收队列大小的代价调低接收队列会限制消费者的潜在吞吐量并且不能用于分区主题partitioned topics。这个「性能 vs 控制」的取舍是否值得完全取决于你的实际使用场景——若吞吐优先保持默认的 1000 甚至更大即可若追求「每个消费者同一时刻只处理一条消息」的强控制则把队列压到 0 或极小值。Java 客户端示例import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.SubscriptionType; String SERVICE_URL pulsar://localhost:6650; String TOPIC persistent://public/default/mq-topic-1; String subscription sub-1; PulsarClient client PulsarClient.builder() .serviceUrl(SERVICE_URL) .build(); Consumer consumer client.newConsumer() .topic(TOPIC) .subscriptionName(subscription) .subscriptionType(SubscriptionType.Shared) // If youd like to restrict the receiver queue size .receiverQueueSize(10) .subscribe();对应的 API 定义位于 ConsumerBuilder.javasubscriptionType(SubscriptionType)与 ConsumerBuilder.javareceiverQueueSize(int)。此外还有maxTotalReceiverQueueSizeAcrossPartitions默认 50000用于限制分区消费者跨分区累计预取的消息总数见 ConsumerBuilder.java。提示Java 消费者还可以通过consumer.receive()同步或receiveAsync()异步拉取消息。配合共享订阅时注意不能使用累积确认cumulative ack只能逐条确认消息——这一点在概念文档中有明确说明。Python 客户端示例from pulsar import Client, ConsumerType SERVICE_URL pulsar://localhost:6650 TOPIC persistent://public/default/mq-topic-1 SUBSCRIPTION sub-1 client Client(SERVICE_URL) consumer client.subscribe( TOPIC, SUBSCRIPTION, # If youd like to restrict the receiver queue size receiver_queue_size10, consumer_typeConsumerType.Shared)在 Python 客户端的init.py 中可以看到receiver_queue_size的默认值同样是1000跨分区累积上限max_total_receiver_queue_size_across_partitions默认50000——与 Java 客户端保持一致。C 客户端示例#include pulsar/Client.h std::string serviceUrl pulsar://localhost:6650; std::string topic persistent://public/default/mq-topic-1; std::string subscription sub-1; Client client(serviceUrl); ConsumerConfiguration consumerConfig; consumerConfig.setConsumerType(ConsumerType.ConsumerShared); // If youd like to restrict the receiver queue size consumerConfig.setReceiverQueueSize(10); Consumer consumer; Result result client.subscribe(topic, subscription, consumerConfig, consumer);C 客户端的相关定义可在 ConsumerConfiguration.hsetReceiverQueueSize与 ConsumerType.hConsumerShared中查看。C API 同样提供pulsar_ConsumerShared枚举见 consumer_configuration.h。Go 客户端示例import github.com/apache/pulsar-client-go/pulsar client, err : pulsar.NewClient(pulsar.ClientOptions{ URL: pulsar://localhost:6650, }) if err ! nil { log.Fatal(err) } consumer, err : client.Subscribe(pulsar.ConsumerOptions{ Topic: persistent://public/default/mq-topic-1, SubscriptionName: sub-1, Type: pulsar.Shared, ReceiverQueueSize: 10, // If youd like to restrict the receiver queue size }) if err ! nil { log.Fatal(err) }原理纵深Shared 订阅与接收队列如何协同工作1. Shared 订阅消息如何分摊给多个消费者在 shared轮询模式下多个消费者可以挂载到同一条订阅。消息以轮询round robin方式在消费者之间分发任意一条消息只会被投递给一个消费者。当某个消费者断开时已投递给它但尚未确认的消息会被重新调度交给剩余消费者处理。从 SubscriptionType.java 的枚举定义可以看到Shared与Key_Shared等类型其中Shared正是消息队列多消费者分摊的核心开关。使用 Shared 类型时需注意两个限制概念文档消息顺序不保证Shared 订阅下消息会在消费者间轮询分发Topic 级别的严格顺序不再保持不能使用累积确认Shared 涉及多个消费者访问同一订阅只能逐条确认消息。2. 接收队列流量控制与精确分派的阀门Pulsar 的流控基于「流许可」flow permit机制消费者向 broker 发送 flow permit 请求来获取消息broker 将消息推送到消费者本地的接收队列。每次调用consumer.receive()消息便从该缓冲中出队。接收队列大小的含义由此清晰起来receiverQueueSize行为表现适用场景默认 1000消费者连接后一次预取最多 1000 条吞吐高但控制粒度粗吞吐优先、对分派公平性要求不高10示例值消费者手头最多预取 10 条控制粒度较细需要较均衡的任务分摊0消费者同一时刻只处理一条消息控制最强每条消息处理代价高、需严格串行化由于接收队列是消费者侧预取缓冲队列越大broker 一次性下发给单个消费者的消息就越多在 shared 订阅下不同消费者实际处理的积压消息数就越不均衡反之队列越小broker 的流控越频繁分派越精确但消息拉取的往返开销上升吞吐随之下降。3. 限制与边界分区主题与吞吐取舍原文档明确指出限制接收队列大小不能与分区主题partitioned topics搭配使用。原因是分区主题被实现为 N 个内部主题跨分区的消息投递与路由语义会使「每个消费者只处理一条消息」的精细控制失效。因此需要精确分派控制时应使用非分区主题 Shared 订阅 小接收队列需要高吞吐时可考虑分区主题 Shared 订阅分区主题下每个分区可挂多个共享消费者但请按场景评估控制粒度需求。从 ConsumerBuilder.java 的注释可知分区场景下消费者跨分区预取总量由maxTotalReceiverQueueSizeAcrossPartitions默认 50000约束这也是为分区主题单独设置流控上限的佐证。实操建议一套代码搭起消息队列结合上文搭建一个「Pulsar 版消息队列」的最小步骤为准备集群本地可先用 standalone 模式启动 broker默认监听pulsar://localhost:6650Topic 无需预先创建客户端首次读写时会自动创建选择主题示例中的persistent://public/default/mq-topic-1即持久化主题消息会写入 BookKeeper保证「处理慢或组件故障时不丢数据」启动消费者集群按上文任意语言示例启动 N 个进程它们必须使用相同的订阅名如sub-1并统一设置SubscriptionType.Shared与目标接收队列大小启动生产者向同一 Topic 发布消息broker 会按轮询策略把消息分派给各消费者一条消息只被一个消费者处理确认与重投消费者处理完每条消息后必须逐条 ack未 ack 的消息会保留在积压中消费者断开时会重新调度给剩余消费者——这正是消息队列「不丢任务」语义的保障。如果需要进一步理解持久化存储如何保证消息不丢失可阅读架构概览若想深入了解消息确认、负确认nack与死信主题等与消息队列配套的可靠性机制可继续阅读消息传递概念文档。小结把 Pulsar Topic 当作消息队列本质上是两种客户端配置的组合Shared 订阅把一条订阅名下的多个消费者变成分摊任务的处理集群小接收队列则把「一次预取 1000 条」的粗粒度流控收紧为「一次只处理少量消息」的精确分派。Java、Python、C、Go 四类客户端 API 语义一致、默认值对齐接收队列默认 1000你可以按团队技术栈自由选择同时要牢记两个边界——Shared 订阅下无顺序保证、不能累积确认小接收队列不能用于分区主题。在「每条消息都必须被处理」的场景下这套配置无需引入额外中间件即可让 Pulsar 同时胜任实时消息总线与消息队列两种角色。【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsar创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考