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;@Slf4jpublic 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;@Slf4jpublic 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
机制进行重试
同步模式举例
这里需要注意,如果处理消息时长差距比较大,同步处理的方式可能会让本来可以很快处理的消息得不到处理的机会。
@Slf4jpublic class DemoMessageListenerSyncAtLeastOnce<T> implements MessageListener<T> {@Overridepublic 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 消息数来阻塞处理。
@Slf4jpublic class DemoMessageListenerAsyncAtLeastOnce<T> implements MessageListener<T> {@Overridepublic void received(Consumer<T> consumer, Message<T> msg) {try {asyncPayload(msg.getData(), new DemoSendCallback() {@Overridepublic 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 控制处理
@Slf4jpublic class DemoMessageListenerBlockListener<T> implements MessageListener<T> {/*** Semaphore保证最多同时处理500条消息*/private final Semaphore semaphore = new Semaphore(500);@Overridepublic void received(Consumer<T> consumer, Message<T> msg) {try {semaphore.acquire();asyncPayload(msg.getData(), new DemoSendCallback() {@Overridepublic 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的其他消息。示例如下:
@Slf4jpublic class DemoMessageListenerSyncAtLeastOnceStrictlyOrdered<T> implements MessageListener<T> {@Overridepublic 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 namespacePattern allTopicsInNamespace = Pattern.compile("public/default/.*");Consumer allTopicsConsumer = consumerBuilder.topicsPattern(allTopicsInNamespace).subscribe();// Subscribe to a subsets of topics in a namespace, based on regexPattern 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 messagereceiveMessageFromConsumer(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 nameString 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 clientPulsarClient pulsarClient = PulsarClient.builder().serviceUrl(localClientUrl).build();//创建consumerConsumer 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)方式。
创建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();
非分区topic支持seek功能
3. peekMessages(topicName,subsciptionName,numMessages)
查询消费者消费到的numMessages条消息。
支持非分区主题(此结论借鉴网上其他博客,还未验证)
不支持分区主题,直接返回null(实测有效)
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);
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 agoconsumer.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();}}}




