暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

pulsar实践(二)

大数据启示录 2022-02-15
871

01

consumer

初始化消费者重要参数

receiverQueueSize

注意:处理不过来时,消费缓冲队列会积压在内存中,合理配置防止 OOM。

autoUpdatePartition

自动更新 partition 信息。如topic中partition信息不变则不需要配置,降低集群的消耗。

subscribeType

订阅类型,根据业务需求决定。

subscriptionInitialPosition

订阅开始的位置,根据业务需求决定最前或者最后。

messageListener

使用 listener 模式消费,只需要提供回调函数,不需要主动执行receive()拉取。一般没有特殊诉求,建议采用 listener 模式。

ackTimeout

当服务端推送消息,但消费者未及时回复 ack 时,经过 ackTimeout 后,会重新推送给消费者处理,即redeliver机制。

注意在利用redeliver机制的时候,一定要注意仅仅使用重试机制来重试可恢复的错误。举个例子,如果代码里面对消息进行解码,解码失败就不适合利用redeliver机制。这会导致客户端一直处于重试之中。

如果拿捏不准,还可以通过下面的deadLetterPolicy配置死信队列,防止消息一直重试。


negativeAckRedeliveryDelay

当客户端调用negativeAcknowledge时,触发redeliver机制的时间。redeliver机制的注意点同ackTimeout。

需要注意的是, ackTimeout和negativeAckRedeliveryDelay建议不要同时使用,一般建议使用negativeAck,用户可以有更灵活的控制权。一旦ackTimeout配置的不合理,在消费时间不确定的情况下可能会导致消息不必要的重试。


deadLetterPolicy

配置redeliver的最大次数和死信 topic。


max total receiver queue size across partitions

设置跨分区中所有接收者队列的最大值。如果超过这个设定值(默认:5000),该设置将会减少各个分区接收者队列的大小(#receiverQueueSize(int))


pattern auto discovery period

该方法也只支持consumer是正则匹配下的那种订阅形式。自动发现的周期以分钟为单位,默认值和最小值为1分钟。

一个消费者一个线程,适用于消费者数目较少的场景

import io.netty.util.concurrent.DefaultThreadFactory;
import lombok.extern.slf4j.Slf4j;
import org.apache.pulsar.client.api.Consumer;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;


@Slf4j
public class DemoPulsarConsumerInit {
private final ScheduledExecutorService executorService = Executors.newScheduledThreadPool(1, new DefaultThreadFactory("pulsar-consumer-init"));
private final String topic;
private volatile Consumer<byte[]> consumer;
public DemoPulsarConsumerInit(String topic) {
this.topic = topic;
}
public void init() {
executorService.scheduleWithFixedDelay(this::initWithRetry, 0, 10, TimeUnit.SECONDS);
}
private void initWithRetry() {
try {
final DemoPulsarClientInit instance = DemoPulsarClientInit.getInstance();
consumer = instance.getPulsarClient().newConsumer().topic(topic).messageListener(new DemoMessageListener<>()).subscribe();
} catch (Exception e) {
log.error("init pulsar producer error, exception is ", e);
}
}
public Consumer<byte[]> getConsumer() {
return consumer;
}
}

多个消费者一个线程,适用于消费者数目较多的场景

import io.netty.util.concurrent.DefaultThreadFactory;
import lombok.extern.slf4j.Slf4j;
import org.apache.pulsar.client.api.Consumer;
import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
@Slf4j
public class DemoPulsarConsumersInit {
private final ScheduledExecutorService executorService = Executors.newScheduledThreadPool(1, new DefaultThreadFactory("pulsar-consumer-init"));
private CopyOnWriteArrayList<Consumer<byte[]>> consumers;
private int initIndex;
private List<String> topics;
public DemoPulsarConsumersInit(List<String> topics) {
this.topics = topics;
}
public void init() {
executorService.scheduleWithFixedDelay(this::initWithRetry, 0, 10, TimeUnit.SECONDS);
}
private void initWithRetry() {
if (initIndex == topics.size()) {
return;
}
for (; initIndex < topics.size(); initIndex++) {
try {
final DemoPulsarClientInit instance = DemoPulsarClientInit.getInstance();
final Consumer<byte[]> consumer = instance.getPulsarClient().newConsumer().topic(topics.get(initIndex)).messageListener(new DemoMessageListener<>()).subscribe();
consumers.add(consumer);
} catch (Exception e) {
log.error("init pulsar producer error, exception is ", e);
break;
}
}
}
public CopyOnWriteArrayList<Consumer<byte[]>> getConsumers() {
return consumers;
}
}

消费者达到至少一次语义

使用手动回复 ack 模式,确保处理成功后再 ack。如果处理失败可以自己重试或通过negativeAck
机制进行重试

同步模式举例

这里需要注意,如果处理消息时长差距比较大,同步处理的方式可能会让本来可以很快处理的消息得不到处理的机会。

@Slf4j
public class DemoMessageListenerSyncAtLeastOnce<T> implements MessageListener<T> {
@Override
public void received(Consumer<T> consumer, Message<T> msg) {
try {
final boolean result = syncPayload(msg.getData());
if (result) {
consumer.acknowledgeAsync(msg);
} else {
consumer.negativeAcknowledge(msg);
}
} catch (Exception e) {
// 业务方法可能会抛出异常
log.error("exception is ", e);
consumer.negativeAcknowledge(msg);
}
}
/**
* 模拟同步执行的业务方法
* @param msg 消息体内容
* @return
*/
private boolean syncPayload(byte[] msg) {
return System.currentTimeMillis() % 2 == 0;
}
}

异步模式举例

异步的话需要考虑内存的限制,因为异步的方式可以很快地从broker消费,不会被业务操作阻塞,这样 inflight 的消息可能会非常多。如果是Shared或KeyShared模式,可以通过maxUnAckedMessage进行限制。如果是Failover模式,可以通过下面的消费者繁忙时阻塞拉取消息,不再进行业务处理通过判断 inflight 消息数来阻塞处理。

@Slf4j
public class DemoMessageListenerAsyncAtLeastOnce<T> implements MessageListener<T> {
@Override
public void received(Consumer<T> consumer, Message<T> msg) {
try {
asyncPayload(msg.getData(), new DemoSendCallback() {
@Override
public void callback(Exception e) {
if (e == null) {
consumer.acknowledgeAsync(msg);
} else {
log.error("exception is ", e);
consumer.negativeAcknowledge(msg);
}
}
});
} catch (Exception e) {
// 业务方法可能会抛出异常
consumer.negativeAcknowledge(msg);
}
}
/**
* 模拟异步执行的业务方法
* @param msg 消息体
* @param demoSendCallback 异步函数的callback
*/
private void asyncPayload(byte[] msg, DemoSendCallback demoSendCallback) {
if (System.currentTimeMillis() % 2 == 0) {
demoSendCallback.callback(null);
} else {
demoSendCallback.callback(new Exception("exception"));
}
}
}

消费者繁忙时阻塞拉取消息,不再进行业务处理

当消费者处理不过来时,通过阻塞listener
方法,不再进行业务处理。避免在微服务积累太多消息导致 OOM,可以通过 RateLimiter 或者 Semaphore 控制处理

@Slf4j
public class DemoMessageListenerBlockListener<T> implements MessageListener<T> {
/**
* Semaphore保证最多同时处理500条消息
*/
private final Semaphore semaphore = new Semaphore(500);
@Override
public void received(Consumer<T> consumer, Message<T> msg) {
try {
semaphore.acquire();
asyncPayload(msg.getData(), new DemoSendCallback() {
@Override
public void callback(Exception e) {
semaphore.release();
if (e == null) {
consumer.acknowledgeAsync(msg);
} else {
log.error("exception is ", e);
consumer.negativeAcknowledge(msg);
}
}
});
} catch (Exception e) {
semaphore.release();
// 业务方法可能会抛出异常
consumer.negativeAcknowledge(msg);
}
}
/**
* 模拟异步执行的业务方法
* @param msg 消息体
* @param demoSendCallback 异步函数的callback
*/
private void asyncPayload(byte[] msg, DemoSendCallback demoSendCallback) {
if (System.currentTimeMillis() % 2 == 0) {
demoSendCallback.callback(null);
} else {
demoSendCallback.callback(new Exception("exception"));
}
}
}

消费者严格按 partition 保序
为了实现partition级别消费者的严格保序,需要对单partition的消息,一旦处理失败,在这条消息重试成功之前不能处理该partition的其他消息。示例如下:

@Slf4j
public class DemoMessageListenerSyncAtLeastOnceStrictlyOrdered<T> implements MessageListener<T> {
@Override
public void received(Consumer<T> consumer, Message<T> msg) {
retryUntilSuccess(msg.getData());
consumer.acknowledgeAsync(msg);
}
private void retryUntilSuccess(byte[] msg) {
while (true) {
try {
final boolean result = syncPayload(msg);
if (result) {
break;
}
} catch (Exception e) {
log.error("exception is ", e);
}
}
}
/**
* 模拟同步执行的业务方法
*
* @param msg 消息体内容
* @return
*/
private boolean syncPayload(byte[] msg) {
return System.currentTimeMillis() % 2 == 0;
}
}

Async receive

Consumer consumer = client.newConsumer()
.topic("my-topic")
.subscriptionName("my-subscription")
.ackTimeout(10, TimeUnit.SECONDS)
.subscriptionType(SubscriptionType.Exclusive)
.subscribe();
 CompletableFuture<Message> asyncMessage = consumer.receiveAsync();       

Batch receive

/**
* 批量消费消息
* 累计消息数量|累计消息大小|指定时间拉取一次满足其一即执行
* @throws PulsarClientException
*/
@GetMapping("/model/comsumerByBatch")
public void comsumerByBatch() throws PulsarClientException {
PulsarClient pulsarFactory = pulsarConf.pulsarFactory();
Consumer<byte[]> consumer = pulsarFactory.newConsumer()
.topic("my-topic")
.subscriptionName("my-subscription")
.batchReceivePolicy(BatchReceivePolicy.builder()
.maxNumMessages(5) //累计消息数量,默认-1,不限制
.maxNumBytes(1024 * 1024)//累计消息大小,默认 10 * 1024 * 1024
.timeout(200, TimeUnit.MILLISECONDS)//指定时间拉取一次,默认 100
.build())
.subscribe();
Messages<byte[]> messages = consumer.batchReceive();
// consumer.batchReceiveAsync();//异步同理
System.out.println("本次接收"+messages.size()+"条");
for (Message<byte[]> message : messages) {
System.out.println(new String(message.getData()));
}
consumer.acknowledge(messages);//确认消息被消费
consumer.close();
    }

negativeAcknowledgement

/**
* 取消消费
* @throws PulsarClientException
*/
@GetMapping("negativeAcknowledgement")
public void negativeAcknowledgement() throws PulsarClientException {
PulsarClient pulsarFactory = pulsarConf.pulsarFactory();
Consumer<byte[]> consumer = pulsarFactory.newConsumer()
.topic("my-topic")
.subscriptionName("my-subscription")
.subscribe();
Message<byte[]> receive = consumer.receive();
MessageId messageId = receive.getMessageId();
try {
System.out.println(new String(receive.getData()));
consumer.acknowledge(receive);//确认消息被消费
} catch (PulsarClientException e) {
consumer.negativeAcknowledge(messageId);//消息取消消费,会被重新消费
e.printStackTrace();
}
consumer.close();
    }

Multi-topic subscriptions

import org.apache.pulsar.client.api.Consumer;
import org.apache.pulsar.client.api.PulsarClient;


import java.util.Arrays;
import java.util.List;
import java.util.regex.Pattern;


ConsumerBuilder consumerBuilder = pulsarClient.newConsumer()
.subscriptionName(subscription);


// Subscribe to all topics in a namespace
Pattern allTopicsInNamespace = Pattern.compile("public/default/.*");
Consumer allTopicsConsumer = consumerBuilder
.topicsPattern(allTopicsInNamespace)
.subscribe();


// Subscribe to a subsets of topics in a namespace, based on regex
Pattern someTopicsInNamespace = Pattern.compile("public/default/foo.*");
Consumer allTopicsConsumer = consumerBuilder
.topicsPattern(someTopicsInNamespace)
.subscribe();




Pattern pattern = Pattern.compile("public/default/.*");
pulsarClient.newConsumer()
.subscriptionName("my-sub")
.topicsPattern(pattern)
.subscriptionTopicsMode(RegexSubscriptionMode.AllTopics)
.subscribe();


Pattern allTopicsInNamespace = Pattern.compile("persistent://public/default.*");
consumerBuilder
.topics(topics)
.subscribeAsync()
.thenAccept(this::receiveMessageFromConsumer);


private void receiveMessageFromConsumer(Object consumer) {
((Consumer)consumer).receiveAsync().thenAccept(message -> {
// Do something with the received message
receiveMessageFromConsumer(consumer);
});
}               

consumer完整demo

import org.apache.pulsar.client.api.Consumer;
import org.apache.pulsar.client.api.Message;
import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.client.api.SubscriptionType;


import java.util.concurrent.TimeUnit;




public class PulsarConsumer {
public static void main(String[] args) throws Exception {
String localClientUrl = "pulsar://xxx.xxx.xx.xx:29095";
// String localClientUrl = "pulsar://10.26.114.120:6650";
// 需要订阅的topic name
String topicName = "persistent://public/default/hsh5-topic";
// 订阅名
String subscriptionName = "my-sub";
consumerPulsarInfo(localClientUrl, topicName, subscriptionName);
}




/**
* 消费数据
*
* @param localClientUrl 消费的主机地址
* @param topicName 主题
* @param subscriptionName 订阅名称
* @throws Exception
*/
public static void consumerPulsarInfo(String localClientUrl, String topicName, String subscriptionName) throws Exception {
// 构造Pulsar client
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(localClientUrl).build();
//创建consumer
Consumer consumer = pulsarClient.newConsumer()
.topic(topicName)
.subscriptionName(subscriptionName)
// 指定消费模式,包含:Exclusive,Failover,Shared,Key_Shared。默认Exclusive模式
.subscriptionType(SubscriptionType.Exclusive)
// 指定从哪里开始消费还有Latest,valueof可选,默认Latest
// .subscriptionInitialPosition(SubscriptionInitialPosition.Earliest)
// 指定消费失败后延迟多久broker重新发送消息给consumer,默认60s
.negativeAckRedeliveryDelay(60, TimeUnit.SECONDS)
.subscribe();
//消费消息
while (true) {
Message message = consumer.receive();
try {
System.out.printf("Message received: %s%n", new String(message.getData()));
consumer.acknowledge(message);
} catch (Exception e) {
e.printStackTrace();
consumer.negativeAcknowledge(message);
}
}
}
}

pulsar重置消费,移动偏移量

Pulsar 消费重置,移动偏移量有6种方法

设置subscriptionInitialPosition,在创建consume的时候处理。

consumer.seek(messageId)方式。

admin.topics().peekMessages(topicName,subsciptionName,numMessages)方式。

admin.topics().resetCursor(topicName,subsciptionName,messageTimestamp)方式。

admin.topics().skipMessages(String topic, String subName, long numMessages)方式。

admin.topics().skipAllMessages(String topic, String subName)方式。

1. 设置subscriptionInitialPosition

创建consume的时候我们可以指定subscriptionInitialPosition(SubscriptionInitialPosition.Earliest)参数;
已经支持的有SubscriptionInitialPosition.Earliest和SubscriptionInitialPosition.Latest,顾名思义,一个是回滚到最初,一个是最新接受到的消息。目标是对某个订阅而言;

PulsarClient client = PulsarClient.builder().serviceUrl(pulsar.getBrokerServiceUrl()).build();


CompletableFuture<Producer<String>> producerFuture = client.newProducer(Schema.STRING)
.topic(topicName)
.createAsync();
CompletableFuture<Consumer<String>> consumerFuture = client.newConsumer(Schema.STRING)
.topic(topicName)
.subscriptionName("sub")
.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest)
.subscribeAsync();


CompletableFuture.allOf(producerFuture, consumerFuture).get();


Producer<String> producer = producerFuture.get();
Consumer<String> consumer = consumerFuture.get();


for (int i = 0; i < N; i++) {
producer.send("Hello-" + i);
}


consumer.close();
producer.close();

2. consumer.seek(messageId)

重置游标
pulsar中有分区topic和非分区topic的区分,现在api只支持创建分区topic,例如我们创建了一个分区topic(persistent://zhiwang3/whds9/admin2),然后创建分区topic的consume,此时我们发现分区topic的consume不支持seek功能。
分区topic不支持seek功能
非分区topic支持seek功能

3. peekMessages(topicName,subsciptionName,numMessages)
查询消费者消费到的numMessages条消息。
支持非分区主题(此结论借鉴网上其他博客,还未验证)
不支持分区主题,直接返回null(实测有效)


4. resetCursor(topicName,subsciptionName,messageTimestamp)

重置游标,
分区topic支持方法(实测有效)
CompletableFuture resetCursorAsync(String topic, String subName, long timestamp);

非分区topic支持方法(实测有效)
void resetCursor(String topic, String subName, MessageId messageId) throws PulsarAdminException;

参数解释:

/**
* Reset cursor position on a topic subscription.
*
* @param topic
* topic name
* @param subName
* Subscription name
* @param messageId
* reset subscription to messageId (or previous nearest messageId if given messageId is not valid)/你可以定义自己的messageId,也可以直接使用已经定义好的MessageId.earliest或MessageId.latest(MessageId.earliest来指向topic上最早
可用的消息,使用MessageId.latest指向最新的消息)
*/
CompletableFuture<Void> resetCursorAsync(String topic, String subName, MessageId messageId);
5. admin.topics().skipMessages(String topic, String subName, long numMessages)

不支持分区topic

6. admin.topics().skipAllMessages(String topic, String subName)

直接跳过所有未被消费的消息。
支持分区和非分区两种topic(实测有效,源码测试代码也可看出)


pulsar client从5小时之前消费:

public class SeekConsumer {
public static String adminUrl = "http://broker1:8080";
public static String serviceUrl = "pulsar://broker2:6650";


public static void main(String[] args) {
try {
PulsarClient client = PulsarClient.builder().serviceUrl(serviceUrl).build();


Consumer<String> consumer = client.newConsumer(Schema.STRING)
.topic("persistent://public/default/test-reset-cursor")
.subscriptionName("seek test")
.subscriptionInitialPosition(SubscriptionInitialPosition.Latest)
.subscribe();


// seek consumer to 5 hours ago
consumer.seek(Instant.now().minus(Duration.ofHours(5)).toEpochMilli());


while (true) {
final Message<String> msg = consumer.receive();


System.out.printf(
"Message received: key=%s, value=%s, topic=%s, id=%s%n",
msg.getKey(),
msg.getValue(),
msg.getTopicName(),
msg.getMessageId().toString());
consumer.acknowledge(msg);
}


} catch (PulsarClientException e) {
e.printStackTrace();
}
}
}




文章转载自大数据启示录,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论