From bf1dbd9547243a0b190a711be98e35eaf8829067 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Tue, 30 Aug 2022 21:06:35 -0400 Subject: [PATCH] Adding test for initial subscription position --- .../core/ConsumerAcknowledgmentTests.java | 36 ++--- .../pulsar/core/DefaultConsumerTests.java | 68 --------- ...tPulsarMessageListenerContainierTests.java | 143 ++++++++++++++++++ 3 files changed, 161 insertions(+), 86 deletions(-) delete mode 100644 spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultConsumerTests.java create mode 100644 spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainierTests.java diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java index 929bab3c..867f21b1 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java @@ -62,9 +62,9 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { void testRecordAck() throws Exception { Map config = new HashMap<>(); final Set strings = new HashSet<>(); - strings.add("foobar-011"); + strings.add("cons-ack-tests-011"); config.put("topicNames", strings); - config.put("subscriptionName", "foobar-sb-011"); + config.put("subscriptionName", "cons-ack-tests-sb-011"); final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build(); final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( pulsarClient, config); @@ -87,7 +87,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { }).when(containerConsumer).acknowledge(any(Message.class)); Map prodConfig = new HashMap<>(); - prodConfig.put("topicName", "foobar-011"); + prodConfig.put("topicName", "cons-ack-tests-011"); final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( pulsarClient, prodConfig); final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); @@ -103,9 +103,9 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { void testBatchAck() throws Exception { Map config = new HashMap<>(); final Set strings = new HashSet<>(); - strings.add("foobar-012"); + strings.add("cons-ack-tests-012"); config.put("topicNames", strings); - config.put("subscriptionName", "foobar-sb-012"); + config.put("subscriptionName", "cons-ack-tests-sb-012"); final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build(); final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( pulsarClient, config); @@ -121,7 +121,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { final Consumer containerConsumer = spyOnConsumer(container); Map prodConfig = new HashMap<>(); - prodConfig.put("topicName", "foobar-012"); + prodConfig.put("topicName", "cons-ack-tests-012"); final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( pulsarClient, prodConfig); final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); @@ -141,9 +141,9 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { void testBatchAckButSomeRecordsFail() throws Exception { Map config = new HashMap<>(); final Set strings = new HashSet<>(); - strings.add("foobar-013"); + strings.add("cons-ack-tests-013"); config.put("topicNames", strings); - config.put("subscriptionName", "foobar-sb-013"); + config.put("subscriptionName", "cons-ack-tests-sb-013"); final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build(); final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( pulsarClient, config); @@ -163,7 +163,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { final Consumer containerConsumer = spyOnConsumer(container); Map prodConfig = new HashMap<>(); - prodConfig.put("topicName", "foobar-013"); + prodConfig.put("topicName", "cons-ack-tests-013"); final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( pulsarClient, prodConfig); final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); @@ -187,9 +187,9 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { void testManualAckForRecordListener() throws Exception { Map config = new HashMap<>(); final Set strings = new HashSet<>(); - strings.add("foobar-014"); + strings.add("cons-ack-tests-014"); config.put("topicNames", strings); - config.put("subscriptionName", "foobar-sb-014"); + config.put("subscriptionName", "cons-ack-tests-sb-014"); final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build(); final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( pulsarClient, config); @@ -218,7 +218,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { }).when(containerConsumer).acknowledge(any(Message.class)); Map prodConfig = new HashMap<>(); - prodConfig.put("topicName", "foobar-014"); + prodConfig.put("topicName", "cons-ack-tests-014"); final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( pulsarClient, prodConfig); final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); @@ -241,9 +241,9 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { void testBatchAckForBatchListener() throws Exception { Map config = new HashMap<>(); final Set strings = new HashSet<>(); - strings.add("foobar-015"); + strings.add("cons-ack-tests-015"); config.put("topicNames", strings); - config.put("subscriptionName", "foobar-sb-015"); + config.put("subscriptionName", "cons-ack-tests-sb-015"); final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build(); final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( pulsarClient, config); @@ -268,7 +268,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { final Consumer containerConsumer = spyOnConsumer(container); Map prodConfig = new HashMap<>(); - prodConfig.put("topicName", "foobar-015"); + prodConfig.put("topicName", "cons-ack-tests-015"); final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( pulsarClient, prodConfig); final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); @@ -289,9 +289,9 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { void testBatchNackForEntireBatchWhenUsingBatchListener() throws Exception { Map config = new HashMap<>(); final Set strings = new HashSet<>(); - strings.add("foobar-016"); + strings.add("cons-ack-tests-016"); config.put("topicNames", strings); - config.put("subscriptionName", "foobar-sb-016"); + config.put("subscriptionName", "cons-ack-tests-sb-016"); final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build(); final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( pulsarClient, config); @@ -316,7 +316,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { final Consumer containerConsumer = spyOnConsumer(container); Map prodConfig = new HashMap<>(); - prodConfig.put("topicName", "foobar-016"); + prodConfig.put("topicName", "cons-ack-tests-016"); final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( pulsarClient, prodConfig); final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultConsumerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultConsumerTests.java deleted file mode 100644 index d91fbbda..00000000 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultConsumerTests.java +++ /dev/null @@ -1,68 +0,0 @@ -/* - * Copyright 2022 the original author or authors. - * - * Licensed 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 - * - * https://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.springframework.pulsar.core; - -import java.util.HashMap; -import java.util.HashSet; -import java.util.Map; -import java.util.concurrent.CompletableFuture; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.TimeUnit; - -import org.apache.pulsar.client.api.MessageId; -import org.apache.pulsar.client.api.PulsarClient; -import org.apache.pulsar.client.api.Schema; -import org.junit.jupiter.api.Test; - -import org.springframework.pulsar.listener.DefaultPulsarMessageListenerContainer; -import org.springframework.pulsar.listener.PulsarContainerProperties; -import org.springframework.pulsar.listener.PulsarRecordMessageListener; - -/** - * @author Soby Chacko - */ -class DefaultConsumerTests extends AbstractContainerBaseTests { - - @Test - void testDefaultConsumer() throws Exception { - Map config = new HashMap<>(); - final HashSet strings = new HashSet(); - strings.add("foobar-012"); - config.put("topicNames", strings); - config.put("subscriptionName", "foobar-sb-012"); - final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( - pulsarClient, config); - CountDownLatch latch = new CountDownLatch(1); - PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); - pulsarContainerProperties - .setMessageListener((PulsarRecordMessageListener) (consumer, msg) -> latch.countDown()); - pulsarContainerProperties.setSchema(Schema.STRING); - DefaultPulsarMessageListenerContainer container = new DefaultPulsarMessageListenerContainer<>( - pulsarConsumerFactory, pulsarContainerProperties); - container.start(); - Map prodConfig = new HashMap<>(); - prodConfig.put("topicName", "foobar-012"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); - final CompletableFuture future = pulsarTemplate.sendAsync("hello john doe"); - latch.await(10, TimeUnit.SECONDS); - pulsarClient.close(); - } - -} diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainierTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainierTests.java new file mode 100644 index 00000000..70593350 --- /dev/null +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainierTests.java @@ -0,0 +1,143 @@ +/* + * Copyright 2022 the original author or authors. + * + * Licensed 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 + * + * https://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.springframework.pulsar.listener; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.util.ArrayList; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; + +import org.apache.pulsar.client.api.PulsarClient; +import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.client.api.SubscriptionInitialPosition; +import org.junit.jupiter.api.Test; + +import org.springframework.pulsar.core.AbstractContainerBaseTests; +import org.springframework.pulsar.core.DefaultPulsarConsumerFactory; +import org.springframework.pulsar.core.DefaultPulsarProducerFactory; +import org.springframework.pulsar.core.PulsarTemplate; + +/** + * @author Soby Chacko + */ +class DefaultPulsarMessageListenerContainierTests extends AbstractContainerBaseTests { + + @Test + void basicDefaultConsumer() throws Exception { + Map config = new HashMap<>(); + final HashSet strings = new HashSet<>(); + strings.add("dpmlct-012"); + config.put("topicNames", strings); + config.put("subscriptionName", "dpmlct-sb-012"); + final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build(); + final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( + pulsarClient, config); + CountDownLatch latch = new CountDownLatch(1); + PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); + pulsarContainerProperties + .setMessageListener((PulsarRecordMessageListener) (consumer, msg) -> latch.countDown()); + pulsarContainerProperties.setSchema(Schema.STRING); + DefaultPulsarMessageListenerContainer container = new DefaultPulsarMessageListenerContainer<>( + pulsarConsumerFactory, pulsarContainerProperties); + container.start(); + Map prodConfig = new HashMap<>(); + prodConfig.put("topicName", "dpmlct-012"); + final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( + pulsarClient, prodConfig); + final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + pulsarTemplate.sendAsync("hello john doe"); + assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); + container.stop(); + pulsarClient.close(); + } + + @Test + void subscriptionInitialPositionEarliest() throws Exception { + Map config = new HashMap<>(); + final HashSet strings = new HashSet<>(); + strings.add("dpmlct-013"); + config.put("topicNames", strings); + config.put("subscriptionName", "dpmlct-sb-013"); + config.put("subscriptionInitialPosition", SubscriptionInitialPosition.Earliest); + final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build(); + final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( + pulsarClient, config); + CountDownLatch latch = new CountDownLatch(5); + PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); + pulsarContainerProperties + .setMessageListener((PulsarRecordMessageListener) (consumer, msg) -> latch.countDown()); + pulsarContainerProperties.setSchema(Schema.STRING); + DefaultPulsarMessageListenerContainer container = new DefaultPulsarMessageListenerContainer<>( + pulsarConsumerFactory, pulsarContainerProperties); + + Map prodConfig = new HashMap<>(); + prodConfig.put("topicName", "dpmlct-013"); + final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( + pulsarClient, prodConfig); + final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + for (int i = 0; i < 5; i++) { + pulsarTemplate.send("hello john doe" + i); + } + // Only start container after all the messages are sent + container.start(); + assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); + container.stop(); + pulsarClient.close(); + } + + @Test + void subscriptionInitialPositionDefaultLatest() throws Exception { + Map config = new HashMap<>(); + final HashSet strings = new HashSet<>(); + strings.add("dpmlct-014"); + config.put("topicNames", strings); + config.put("subscriptionName", "dpmlct-sb-014"); + final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build(); + final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( + pulsarClient, config); + PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); + List messages = new ArrayList<>(); + pulsarContainerProperties.setMessageListener( + (PulsarRecordMessageListener) (consumer, msg) -> messages.add((String) msg.getValue())); + pulsarContainerProperties.setSchema(Schema.STRING); + DefaultPulsarMessageListenerContainer container = new DefaultPulsarMessageListenerContainer<>( + pulsarConsumerFactory, pulsarContainerProperties); + + Map prodConfig = new HashMap<>(); + prodConfig.put("topicName", "dpmlct-014"); + final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( + pulsarClient, prodConfig); + final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + for (int i = 0; i < 5; i++) { + pulsarTemplate.send("hello john doe" + i); + } + // Only start container after all the messages are sent + container.start(); + pulsarTemplate.send("hello john doe" + 5); + Thread.sleep(2_000); + assertThat(messages.size()).isEqualTo(1); + assertThat(messages.get(0)).isEqualTo("hello john doe5"); + container.stop(); + pulsarClient.close(); + } + +}