- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
消息队列是大型数据架构中的关键组件:当系统中某个组件变慢甚至宕机时,未被处理的数据必须被可靠保留、并按正确顺序等待后续处理。Apache Pulsar 天生适合承担消息队列的角色——它内置持久化消息存储,并能在同一 topic 的多个 consumer 之间自动做负载均衡(也支持自定义负载均衡)。本文基于 Pulsar 官方 cookbook 文档,深入讲解如何通过"共享订阅 + 控制 receiver queue 大小"把 Pulsar topic 变成标准消息队列,并给出 Java、Python、C++、Go 四种客户端的完整可运行示例与源码级原理佐证。读完本文,你将掌握消息队列场景下的订阅模型选型、receiver queue 调优原理,以及多语言客户端的落地写法。
同一套 Pulsar 集群既可以充当实时消息总线,也可以充当消息队列(或两者兼用)。你可以把一部分 topic 用于实时流式处理,另一部分 topic 用于消息队列场景;也可以为不同用途划分不同 namespace。
为什么 Pulsar 适合做消息队列
Pulsar 的架构天然满足消息队列的两大核心诉求:
- 持久化消息存储:Pulsar 使用 Apache BookKeeper 作为持久化消息存储层。BookKeeper 是分布式预写日志(write-ahead log)系统,为 Pulsar 提供低延迟持久化、跨 bookie 的弹性存储扩展、以及跨数据中心的高可用复制能力。这也是所有持久化 topic 名称中带
persistent前缀的原因(见 concepts-architecture-overview.md)。消息一旦写入即可持久保留,即便消费端组件缓慢或失败,未确认(unacknowledged)的消息也不会丢失。 - 消费者负载均衡:通过共享订阅(Shared subscription),同一订阅名下可挂多个消费者,broker 以 round robin 方式把消息分发给各个消费者,每条消息只投递给一个消费者(详见下文)。
把 Pulsar topic 变成消息队列的两个关键配置
要把 Pulsar topic 用作消息队列,需要从"点对点投递"切换到"工作队列分发",核心是两件事:
建立共享订阅(shared subscription),并让所有消费者使用同一个订阅名。 如果每个消费者使用不同的订阅名,订阅便无法共享,消费者之间也就无法组成一个协作处理整体。共享订阅下,多个消费者可以挂到同一个订阅上,消息以 round robin 方式在消费者之间分发,任意一条消息只会投递给一个消费者;当某个消费者断开时,已发送但未确认的消息会被重新调度给其余消费者(参见 concepts-messaging.md 的 Shared 订阅说明)。
如果需要严格控制消息在消费者之间的分发,把 receiver queue 设得很小(必要时可设为 0)。 每个 Pulsar 消费者都有一个 receiver queue,它决定消费者一次会尝试预取多少条消息。例如默认的 1000 意味着消费者一连接就会尝试从 topic 积压中取出 1000 条消息处理。把 receiver queue 设为 0,本质上就是确保每个消费者同一时刻只处理一件事。
在 Java 客户端 API 的文档注释中(ConsumerBuilder.java),这一行为被描述得更精确:
将消费者队列大小设为 0 会降低消费者吞吐量(因为禁用了消息预取),但能改善共享订阅下的消息分发——broker 只把消息推送给"当前空闲、准备好处理"的消费者。设为 0 时不能使用
receive(int, TimeUnit),也不能使用分区 topic;receive()调用不应被打断。同时不支持批量消息(batch message):若消费者收到批量消息,会关闭与 broker 的连接,receive()将保持阻塞,receiveAsync()则在回调中收到异常,只有排空管线中的批量消息后才能继续接收。
从客户端实现看,ConsumerImpl.java 中receiverQueueRefillThreshold被初始化为conf.getReceiverQueueSize() / 2,即当队列中可用消息降到队列大小的一半时触发补货(refill)请求。queue size 越小,broker 单次推送越少,分发颗粒度越细,消费者的空闲信号越快反馈给 broker。
限制 receiver queue 的代价
限制 receiver queue 的代价是:
- 限制消费者的潜在吞吐量:预取数量减少,消费者"忙等"概率降低,但单消费者并发处理能力也随之下降。
- 不能用于分区 topic(partitioned topics):Java API 文档明确说明 queue size 为 0 时不能与分区 topic 搭配使用(ConsumerBuilder.java)。
吞吐/控制之间的取舍是否值得,取决于你的具体场景:消息处理昂贵、需要严格公平分发时,小 queue size 值得;处理快、追求高吞吐时,保持默认即可。
提示:Pulsar 订阅模型非常灵活(concepts-messaging.md):每个消费者使用唯一订阅名(Exclusive)可实现传统"扇出 pub-sub";多个消费者共享同一订阅名(Shared/Failover/Key_Shared)可实现"消息队列";还可以把两种订阅组合起来,同一 topic 上同时获得 pub-sub 与队列两种效果。
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 you'd like to restrict the receiver queue size .receiverQueueSize(10) .subscribe();要点:
subscriptionName("sub-1")必须与其他消费者保持一致,才能构成共享订阅。SubscriptionType.Shared是枚举值,与Failover、Key_Shared并列定义于 SubscriptionType.java。receiverQueueSize(10)将默认的 1000 降为 10;如要追求极致的单消费者单任务,可传入 0(注意分区 topic 限制)。Java 还提供了maxTotalReceiverQueueSizeAcrossPartitions,用于跨分区限制消费者被 broker 一次推送的消息总数上限(默认 50000,见 ConsumerBuilder.java)。
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 you'd like to restrict the receiver queue size receiver_queue_size=10, consumer_type=ConsumerType.Shared)要点:consumer_type=ConsumerType.Shared对应 Java 的SubscriptionType.Shared;receiver_queue_size=10对应 Java 的receiverQueueSize(10),语义完全一致。
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 you'd like to restrict the receiver queue size consumerConfig.setReceiverQueueSize(10); Consumer consumer; Result result = client.subscribe(topic, subscription, consumerConfig, consumer);要点:C++ 客户端通过ConsumerConfiguration配置ConsumerType.ConsumerShared与setReceiverQueueSize(10),再传给client.subscribe。与消息队列相关的ConsumerType、ConsumerConfiguration等定义位于 pulsar-client-cpp/include/pulsar 目录的 C++ 头文件中。
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 you'd like to restrict the receiver queue size }) if err != nil { log.Fatal(err) }要点:Go 客户端在ConsumerOptions中设置Type: pulsar.Shared与ReceiverQueueSize: 10。Go 客户端是独立仓库(pulsar-client-go),与 Java/Python/C++ 客户端同属 Pulsar 官方客户端体系。
消费端行为与消息队列语义的印证
共享订阅的消息队列语义在客户端实现中有直接体现:
- 分发机制:Shared 类型下,多条消息不会保证全局顺序,且不能使用累积确认(cumulative acknowledgment),只能逐条确认(见 concepts-messaging.md 的限制说明与 concepts-messaging.md 关于负确认的讨论)。
- 单条重投递:在 ConsumerImpl.java 中可以看到,只有 Shared 和 Key_Shared 订阅类型支持对单条消息进行重投递(redelivery)——这正是消息队列"失败重试、换人处理"的基础能力。
这两点共同构成了消息队列的经典语义:多条消费者并发取任务、逐条确认、失败消息重投给其他消费者。
结语
把 Pulsar 当作消息队列,本质上只做两件事:让多个消费者共享同一个订阅名(Shared 订阅)实现负载均衡,以及按需收紧 receiver queue实现精细分发控制。你可以沿用默认的 1000 追求吞吐,也可以压到 0 换取严格的"一人一事";同一集群里还能混合使用不同订阅类型,让实时总线和消息队列共存于一套部署之中。结合本文的多语言示例与源码佐证,你已经可以在自己的项目里快速落地这一模式。
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
使用 Apache Pulsar 构建消息队列(Message Queue):Shared 订阅与 Receiver Queue 实战指南
使用 Apache Pulsar 构建消息队列(Message Queue):Shared 订阅与 Receiver Queue 实战指南 导读 本文基于 Ap
消息队列后端流处理Apache Pulsar 消息队列实践:通过 Shared 订阅与 Receiver Queue 将 Topic 用作消息队列
Apache Pulsar 消息队列实践:通过 Shared 订阅与 Receiver Queue 将 Topic 用作消息队列 导读 :本文围绕 Apache
消息队列后端流处理Apache Pulsar 消息队列实战指南:基于 Shared 订阅与 Receiver Queue 构建可扩展的 MQ 工作负载
Apache Pulsar 消息队列实战指南:基于 Shared 订阅与 Receiver Queue 构建可扩展的 MQ 工作负载 Pulsar 天生具备消息
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考