From d4ef57593a3f3113bb134c4532adddc1d5a81ac3 Mon Sep 17 00:00:00 2001 From: Jia Zhai Date: Mon, 17 Feb 2020 10:35:50 +0800 Subject: [PATCH 1/2] make sub-mode public, added it into ConsumerConfigurationData --- .../pulsar/client/impl/RawReaderImpl.java | 1 - .../pulsar/client/api/SubscriptionMode.java | 31 ++++++++++++++++++ .../pulsar/client/impl/ConsumerImpl.java | 32 +++++++++---------- .../client/impl/MultiTopicsConsumerImpl.java | 13 ++++---- .../pulsar/client/impl/PulsarClientImpl.java | 3 +- .../apache/pulsar/client/impl/ReaderImpl.java | 4 +-- .../client/impl/ZeroQueueConsumerImpl.java | 4 +-- .../impl/conf/ConsumerConfigurationData.java | 3 ++ .../pulsar/client/impl/ConsumerImplTest.java | 27 ++++++++++------ 9 files changed, 78 insertions(+), 40 deletions(-) create mode 100644 pulsar-client-api/src/main/java/org/apache/pulsar/client/api/SubscriptionMode.java diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/client/impl/RawReaderImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/client/impl/RawReaderImpl.java index 6a4d682d2b8ed..3ca072dda51aa 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/client/impl/RawReaderImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/client/impl/RawReaderImpl.java @@ -117,7 +117,6 @@ static class RawConsumerImpl extends ConsumerImpl { TopicName.getPartitionIndex(conf.getSingleTopic()), false, consumerFuture, - SubscriptionMode.Durable, MessageId.earliest, 0 /* startMessageRollbackDurationInSec */, Schema.BYTES, null, diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/SubscriptionMode.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/SubscriptionMode.java new file mode 100644 index 0000000000000..7e11bbb0dc7eb --- /dev/null +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/SubscriptionMode.java @@ -0,0 +1,31 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.client.api; + +/** + * Types of subscription mode supported by Pulsar. + */ +public enum SubscriptionMode { + // Make the subscription to be backed by a durable cursor that will retain messages and persist the current + // position + Durable, + + // Lightweight subscription mode that doesn't have a durable cursor associated + NonDurable +} diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java index 919a280238676..5b62248370c54 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java @@ -63,6 +63,7 @@ import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionInitialPosition; +import org.apache.pulsar.client.api.SubscriptionMode; import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.client.api.PulsarClientException.TopicDoesNotExistException; import org.apache.pulsar.client.impl.conf.ConsumerConfigurationData; @@ -150,39 +151,38 @@ public class ConsumerImpl extends ConsumerBase implements ConnectionHandle private final boolean createTopicIfDoesNotExist; - enum SubscriptionMode { - // Make the subscription to be backed by a durable cursor that will retain messages and persist the current - // position - Durable, - // Lightweight subscription mode that doesn't have a durable cursor associated - NonDurable - } - - static ConsumerImpl newConsumerImpl(PulsarClientImpl client, String topic, ConsumerConfigurationData conf, - ExecutorService listenerExecutor, int partitionIndex, boolean hasParentConsumer, CompletableFuture> subscribeFuture, - SubscriptionMode subscriptionMode, MessageId startMessageId, Schema schema, ConsumerInterceptors interceptors, - boolean createTopicIfDoesNotExist) { + static ConsumerImpl newConsumerImpl(PulsarClientImpl client, + String topic, + ConsumerConfigurationData conf, + ExecutorService listenerExecutor, + int partitionIndex, + boolean hasParentConsumer, + CompletableFuture> subscribeFuture, + MessageId startMessageId, + Schema schema, + ConsumerInterceptors interceptors, + boolean createTopicIfDoesNotExist) { if (conf.getReceiverQueueSize() == 0) { return new ZeroQueueConsumerImpl<>(client, topic, conf, listenerExecutor, partitionIndex, hasParentConsumer, subscribeFuture, - subscriptionMode, startMessageId, schema, interceptors, + startMessageId, schema, interceptors, createTopicIfDoesNotExist); } else { return new ConsumerImpl<>(client, topic, conf, listenerExecutor, partitionIndex, hasParentConsumer, - subscribeFuture, subscriptionMode, startMessageId, 0 /* rollback time in sec to start msgId */, + subscribeFuture, startMessageId, 0 /* rollback time in sec to start msgId */, schema, interceptors, createTopicIfDoesNotExist); } } protected ConsumerImpl(PulsarClientImpl client, String topic, ConsumerConfigurationData conf, ExecutorService listenerExecutor, int partitionIndex, boolean hasParentConsumer, - CompletableFuture> subscribeFuture, SubscriptionMode subscriptionMode, MessageId startMessageId, + CompletableFuture> subscribeFuture, MessageId startMessageId, long startMessageRollbackDurationInSec, Schema schema, ConsumerInterceptors interceptors, boolean createTopicIfDoesNotExist) { super(client, topic, conf, conf.getReceiverQueueSize(), listenerExecutor, subscribeFuture, schema, interceptors); this.consumerId = client.newConsumerId(); - this.subscriptionMode = subscriptionMode; + this.subscriptionMode = conf.getSubscriptionMode(); this.startMessageId = startMessageId != null ? new BatchMessageIdImpl((MessageIdImpl) startMessageId) : null; this.lastDequeuedMessage = startMessageId == null ? MessageId.earliest : startMessageId; this.initialStartMessageId = this.startMessageId; diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java index 9c95d17d9ea4e..3e9e69ae51c3c 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java @@ -59,7 +59,6 @@ import org.apache.pulsar.client.api.PulsarClientException.NotSupportedException; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionType; -import org.apache.pulsar.client.impl.ConsumerImpl.SubscriptionMode; import org.apache.pulsar.client.impl.conf.ConsumerConfigurationData; import org.apache.pulsar.client.impl.transaction.TransactionImpl; import org.apache.pulsar.client.util.ConsumerName; @@ -833,10 +832,10 @@ private void doSubscribeTopicPartitions(CompletableFuture subscribeResult, String partitionName = TopicName.get(topicName).getPartition(partitionIndex).toString(); CompletableFuture> subFuture = new CompletableFuture<>(); ConsumerImpl newConsumer = ConsumerImpl.newConsumerImpl(client, partitionName, - configurationData, client.externalExecutorProvider().getExecutor(), - partitionIndex, true, subFuture, - SubscriptionMode.Durable, null, schema, interceptors, - createIfDoesNotExist); + configurationData, client.externalExecutorProvider().getExecutor(), + partitionIndex, true, subFuture, + null, schema, interceptors, + createIfDoesNotExist); consumers.putIfAbsent(newConsumer.getTopic(), newConsumer); return subFuture; }) @@ -847,7 +846,7 @@ private void doSubscribeTopicPartitions(CompletableFuture subscribeResult, CompletableFuture> subFuture = new CompletableFuture<>(); ConsumerImpl newConsumer = ConsumerImpl.newConsumerImpl(client, topicName, internalConfig, - client.externalExecutorProvider().getExecutor(), -1, true, subFuture, SubscriptionMode.Durable, null, + client.externalExecutorProvider().getExecutor(), -1, true, subFuture, null, schema, interceptors, createIfDoesNotExist); consumers.putIfAbsent(newConsumer.getTopic(), newConsumer); @@ -1118,7 +1117,7 @@ private CompletableFuture subscribeIncreasedTopicPartitions(String topicNa ConsumerImpl newConsumer = ConsumerImpl.newConsumerImpl( client, partitionName, configurationData, client.externalExecutorProvider().getExecutor(), - partitionIndex, true, subFuture, SubscriptionMode.Durable, null, schema, interceptors, + partitionIndex, true, subFuture, null, schema, interceptors, true /* createTopicIfDoesNotExist */); consumers.putIfAbsent(newConsumer.getTopic(), newConsumer); if (log.isDebugEnabled()) { diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java index 5a2a6dae09f60..f51fb6b65c3f3 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java @@ -64,7 +64,6 @@ import org.apache.pulsar.client.api.schema.SchemaInfoProvider; import org.apache.pulsar.client.api.AuthenticationFactory; import org.apache.pulsar.client.api.transaction.TransactionBuilder; -import org.apache.pulsar.client.impl.ConsumerImpl.SubscriptionMode; import org.apache.pulsar.client.impl.conf.ClientConfigurationData; import org.apache.pulsar.client.impl.conf.ConsumerConfigurationData; import org.apache.pulsar.client.impl.conf.ProducerConfigurationData; @@ -354,7 +353,7 @@ private CompletableFuture> doSingleTopicSubscribeAsync(ConsumerC } else { int partitionIndex = TopicName.getPartitionIndex(topic); consumer = ConsumerImpl.newConsumerImpl(PulsarClientImpl.this, topic, conf, listenerThread, partitionIndex, false, - consumerSubscribedFuture, SubscriptionMode.Durable, null, schema, interceptors, + consumerSubscribedFuture,null, schema, interceptors, true /* createTopicIfDoesNotExist */); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ReaderImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ReaderImpl.java index 5022e11fc0814..4f172e39c8401 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ReaderImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ReaderImpl.java @@ -26,7 +26,6 @@ import org.apache.commons.codec.digest.DigestUtils; import org.apache.commons.lang3.StringUtils; import org.apache.pulsar.client.api.*; -import org.apache.pulsar.client.impl.ConsumerImpl.SubscriptionMode; import org.apache.pulsar.client.impl.conf.ConsumerConfigurationData; import org.apache.pulsar.client.impl.conf.ReaderConfigurationData; import org.apache.pulsar.common.naming.TopicName; @@ -47,6 +46,7 @@ public ReaderImpl(PulsarClientImpl client, ReaderConfigurationData readerConf consumerConfiguration.getTopicNames().add(readerConfiguration.getTopicName()); consumerConfiguration.setSubscriptionName(subscription); consumerConfiguration.setSubscriptionType(SubscriptionType.Exclusive); + consumerConfiguration.setSubscriptionMode(SubscriptionMode.NonDurable); consumerConfiguration.setReceiverQueueSize(readerConfiguration.getReceiverQueueSize()); consumerConfiguration.setReadCompacted(readerConfiguration.isReadCompacted()); @@ -91,7 +91,7 @@ public void reachedEndOfTopic(Consumer consumer) { final int partitionIdx = TopicName.getPartitionIndex(readerConfiguration.getTopicName()); consumer = new ConsumerImpl<>(client, readerConfiguration.getTopicName(), consumerConfiguration, - listenerExecutor, partitionIdx, false, consumerFuture, SubscriptionMode.NonDurable, + listenerExecutor, partitionIdx, false, consumerFuture, readerConfiguration.getStartMessageId(), readerConfiguration.getStartMessageFromRollbackDurationInSec(), schema, null, true /* createTopicIfDoesNotExist */); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ZeroQueueConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ZeroQueueConsumerImpl.java index b9249828bc1b1..94c8dd3c56ea4 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ZeroQueueConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ZeroQueueConsumerImpl.java @@ -49,11 +49,11 @@ public class ZeroQueueConsumerImpl extends ConsumerImpl { public ZeroQueueConsumerImpl(PulsarClientImpl client, String topic, ConsumerConfigurationData conf, ExecutorService listenerExecutor, int partitionIndex, boolean hasParentConsumer, CompletableFuture> subscribeFuture, - SubscriptionMode subscriptionMode, MessageId startMessageId, Schema schema, + MessageId startMessageId, Schema schema, ConsumerInterceptors interceptors, boolean createTopicIfDoesNotExist) { super(client, topic, conf, listenerExecutor, partitionIndex, hasParentConsumer, subscribeFuture, - subscriptionMode, startMessageId, 0 /* startMessageRollbackDurationInSec */, schema, interceptors, + startMessageId, 0 /* startMessageRollbackDurationInSec */, schema, interceptors, createTopicIfDoesNotExist); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ConsumerConfigurationData.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ConsumerConfigurationData.java index 6195b5f720844..e64db7287fea8 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ConsumerConfigurationData.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ConsumerConfigurationData.java @@ -43,6 +43,7 @@ import org.apache.pulsar.client.api.MessageListener; import org.apache.pulsar.client.api.RegexSubscriptionMode; import org.apache.pulsar.client.api.SubscriptionInitialPosition; +import org.apache.pulsar.client.api.SubscriptionMode; import org.apache.pulsar.client.api.SubscriptionType; @Data @@ -59,6 +60,8 @@ public class ConsumerConfigurationData implements Serializable, Cloneable { private SubscriptionType subscriptionType = SubscriptionType.Exclusive; + private SubscriptionMode subscriptionMode = SubscriptionMode.Durable; + @JsonIgnore private MessageListener messageListener; diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java index a75765164f73c..dcb2e36e2ef22 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java @@ -18,25 +18,32 @@ */ package org.apache.pulsar.client.impl; +import static org.mockito.Mockito.any; +import static org.mockito.Mockito.anyString; +import static org.mockito.Mockito.doNothing; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + import io.netty.util.Timer; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; + import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.PulsarClientException; -import org.apache.pulsar.client.impl.ConsumerImpl.SubscriptionMode; import org.apache.pulsar.client.impl.conf.ClientConfigurationData; import org.apache.pulsar.client.impl.conf.ConsumerConfigurationData; import org.testng.Assert; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Test; -import java.util.concurrent.CompletableFuture; -import java.util.concurrent.CompletionException; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.TimeUnit; - -import static org.mockito.Mockito.*; - public class ConsumerImplTest { @@ -62,7 +69,7 @@ public void setUp() { consumerConf.setSubscriptionName("test-sub"); consumer = ConsumerImpl.newConsumerImpl(client, topic, consumerConf, - executorService, -1, false, subscribeFuture, SubscriptionMode.Durable, null, null, null, + executorService, -1, false, subscribeFuture, null, null, null, true); } From 6a96226b4e4ee6117563c2e5cf9f636b36ab722d Mon Sep 17 00:00:00 2001 From: Jia Zhai Date: Mon, 17 Feb 2020 10:51:05 +0800 Subject: [PATCH 2/2] add sub mode config in consumerBuilder --- .../apache/pulsar/client/api/ConsumerBuilder.java | 15 +++++++++++++++ .../pulsar/client/impl/ConsumerBuilderImpl.java | 8 ++++++++ 2 files changed, 23 insertions(+) diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java index 50328ff04a392..102d485f54be3 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java @@ -228,6 +228,21 @@ public interface ConsumerBuilder extends Cloneable { */ ConsumerBuilder subscriptionType(SubscriptionType subscriptionType); + /** + * Select the subscription mode to be used when subscribing to the topic. + * + *

Options are: + *

    + *
  • {@link SubscriptionMode#Durable} (Default)
  • + *
  • {@link SubscriptionMode#NonDurable}
  • + *
+ * + * @param subscriptionMode + * the subscription mode value + * @return the consumer builder instance + */ + ConsumerBuilder subscriptionMode(SubscriptionMode subscriptionMode); + /** * Sets a {@link MessageListener} for the consumer * diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBuilderImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBuilderImpl.java index 872abb7b24f43..e31cc84ec50b7 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBuilderImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBuilderImpl.java @@ -47,6 +47,7 @@ import org.apache.pulsar.client.api.PulsarClientException.InvalidConfigurationException; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionInitialPosition; +import org.apache.pulsar.client.api.SubscriptionMode; import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.client.impl.conf.ConfigurationDataUtils; import org.apache.pulsar.client.impl.conf.ConsumerConfigurationData; @@ -191,6 +192,13 @@ public ConsumerBuilder subscriptionType(@NonNull SubscriptionType subscriptio return this; } + @Override + public ConsumerBuilder subscriptionMode(@NonNull SubscriptionMode subscriptionMode) { + conf.setSubscriptionMode(subscriptionMode); + return this; + } + + @Override public ConsumerBuilder messageListener(@NonNull MessageListener messageListener) { conf.setMessageListener(messageListener);