Using Pulsar as a message queue
消息队列是许多大规模数据架构的基本组件。 If every single work object that passes through your system absolutely must be processed in spite of the slowness or downright failure of this or that system component, there’s a good chance that you’ll need a message queue to step in and ensure that unprocessed data is retained—-with correct ordering—-until the required actions are taken.
对于消息队列而言,Pulsar 是绝佳选择,这是因为:
你可以使用相同的 Pulsar 安装配置来充当实时消息总线和消息队列(也可只使用其中一个)。 你可以为实时目的留出一些主题,为消息队列目的留出其他主题(也可为任一目的使用特定的命名空间)。
客户端配置更改
要将 Pulsar 主题用作消息队列,应将该主题的接收器负载分配到多个消费者(最佳消费者数取决于实时负载量)。 每个消费者必须:
建立共享订阅,使用与其他消费者相同的订阅名称(否则订阅将不共享,消费者集群无法充当处理集合)
If you’d like to have tight control over message dispatching across consumers, set the consumers’ receiver queue size very low (potentially even to 0 if necessary). 每个 Pulsar 消费者都有一个接收器队列,用于确定消费者一次尝试获取的消息数量。 例如,接收器队列 1000 (默认值)意味着消费者将尝试在连接时处理来自主题的 1000 条待办消息。 将接收器队列值设置为零实质上意味着确保每个消费者一次只做一件事。
限制消费者的接收器队列大小的缺点是限制了这些消费者的潜在吞吐量,并且不能与分区主题一起使用。 降低性能以获得更好的控制是否值得取决于你的用例。
Java 客户端
以下是使用共享订阅的 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)
// 如果要限制接收器队列大小
.receiverQueueSize(10)
.subscribe();
Python 客户端
以下是使用共享订阅的 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,
# 如果要限制接收器队列大小
receiver_queue_size=10,
consumer_type=ConsumerType.Shared)
C++ 客户端
以下是使用共享订阅的 C++ 消费者配置示例:
#include <pulsar/Client.h>
std::string serviceUrl = "pulsar://localhost:6650";
std::string topic = "persistent://public/defaultmq-topic-1";
std::string subscription = "sub-1";
Client client(serviceUrl);
ConsumerConfiguration consumerConfig;
consumerConfig.setConsumerType(ConsumerType.ConsumerShared);
// 如果要限制接收器队列大小
consumerConfig.setReceiverQueueSize(10);
Consumer consumer;
Result result = client.subscribe(topic, subscription, consumerConfig, consumer);
Go 客户端
以下是使用共享订阅的 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)
}