拓冰建站拓冰建站
首页 / 资讯中心 / 正文

用 Apache Pulsar 构建消息队列:共享订阅、receiver queue 与多语言客户端配置实战

消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载消息队列是大型数据架构中的关键组件当系统中某个组件变慢甚至宕机时未被处理的数据必须被可靠保留、并按正确顺序等待后续处理。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)也不能使用分区 topicreceive()调用不应被打断。同时不支持批量消息batch message若消费者收到批量消息会关闭与 broker 的连接receive()将保持阻塞receiveAsync()则在回调中收到异常只有排空管线中的批量消息后才能继续接收。从客户端实现看ConsumerImpl.java 中receiverQueueRefillThreshold被初始化为conf.getReceiverQueueSize() / 2即当队列中可用消息降到队列大小的一半时触发补货refill请求。queue size 越小broker 单次推送越少分发颗粒度越细消费者的空闲信号越快反馈给 broker。限制 receiver queue 的代价限制 receiver queue 的代价是限制消费者的潜在吞吐量预取数量减少消费者忙等概率降低但单消费者并发处理能力也随之下降。不能用于分区 topicpartitioned topicsJava 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 youd 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 youd like to restrict the receiver queue size receiver_queue_size10, consumer_typeConsumerType.Shared)要点consumer_typeConsumerType.Shared对应 Java 的SubscriptionType.Sharedreceiver_queue_size10对应 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 youd 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 youd 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 换取严格的一人一事同一集群里还能混合使用不同订阅类型让实时总线和消息队列共存于一套部署之中。结合本文的多语言示例与源码佐证你已经可以在自己的项目里快速落地这一模式。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐使用 Apache Pulsar 构建消息队列Message QueueShared 订阅与 Receiver Queue 实战指南使用 Apache Pulsar 构建消息队列Message QueueShared 订阅与 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 天生具备消息消息队列后端流处理上一篇在Android设备上运行完整Linux系统的终极解决方案PRoot-Distro深度指南下一篇MATLAB机器人工具箱终极指南从入门到精通的免费开源解决方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门