Adding test for initial subscription position

This commit is contained in:
Soby Chacko
2022-08-30 21:06:35 -04:00
parent 5ffcf5dc06
commit bf1dbd9547
3 changed files with 161 additions and 86 deletions

View File

@@ -62,9 +62,9 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
void testRecordAck() throws Exception {
Map<String, Object> config = new HashMap<>();
final Set<String> 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<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
pulsarClient, config);
@@ -87,7 +87,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
}).when(containerConsumer).acknowledge(any(Message.class));
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "foobar-011");
prodConfig.put("topicName", "cons-ack-tests-011");
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
pulsarClient, prodConfig);
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -103,9 +103,9 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
void testBatchAck() throws Exception {
Map<String, Object> config = new HashMap<>();
final Set<String> 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<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
pulsarClient, config);
@@ -121,7 +121,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
final Consumer<?> containerConsumer = spyOnConsumer(container);
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "foobar-012");
prodConfig.put("topicName", "cons-ack-tests-012");
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
pulsarClient, prodConfig);
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -141,9 +141,9 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
void testBatchAckButSomeRecordsFail() throws Exception {
Map<String, Object> config = new HashMap<>();
final Set<String> 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<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
pulsarClient, config);
@@ -163,7 +163,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
final Consumer<?> containerConsumer = spyOnConsumer(container);
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "foobar-013");
prodConfig.put("topicName", "cons-ack-tests-013");
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
pulsarClient, prodConfig);
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -187,9 +187,9 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
void testManualAckForRecordListener() throws Exception {
Map<String, Object> config = new HashMap<>();
final Set<String> 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<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
pulsarClient, config);
@@ -218,7 +218,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
}).when(containerConsumer).acknowledge(any(Message.class));
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "foobar-014");
prodConfig.put("topicName", "cons-ack-tests-014");
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
pulsarClient, prodConfig);
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -241,9 +241,9 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
void testBatchAckForBatchListener() throws Exception {
Map<String, Object> config = new HashMap<>();
final Set<String> 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<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
pulsarClient, config);
@@ -268,7 +268,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
final Consumer<?> containerConsumer = spyOnConsumer(container);
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "foobar-015");
prodConfig.put("topicName", "cons-ack-tests-015");
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
pulsarClient, prodConfig);
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -289,9 +289,9 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
void testBatchNackForEntireBatchWhenUsingBatchListener() throws Exception {
Map<String, Object> config = new HashMap<>();
final Set<String> 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<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
pulsarClient, config);
@@ -316,7 +316,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
final Consumer<?> containerConsumer = spyOnConsumer(container);
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "foobar-016");
prodConfig.put("topicName", "cons-ack-tests-016");
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
pulsarClient, prodConfig);
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);

View File

@@ -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<String, Object> config = new HashMap<>();
final HashSet<String> strings = new HashSet<String>();
strings.add("foobar-012");
config.put("topicNames", strings);
config.put("subscriptionName", "foobar-sb-012");
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
final DefaultPulsarConsumerFactory<String> 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<String> container = new DefaultPulsarMessageListenerContainer<>(
pulsarConsumerFactory, pulsarContainerProperties);
container.start();
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "foobar-012");
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
pulsarClient, prodConfig);
final PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
final CompletableFuture<MessageId> future = pulsarTemplate.sendAsync("hello john doe");
latch.await(10, TimeUnit.SECONDS);
pulsarClient.close();
}
}

View File

@@ -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<String, Object> config = new HashMap<>();
final HashSet<String> 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<String> 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<String> container = new DefaultPulsarMessageListenerContainer<>(
pulsarConsumerFactory, pulsarContainerProperties);
container.start();
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "dpmlct-012");
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
pulsarClient, prodConfig);
final PulsarTemplate<String> 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<String, Object> config = new HashMap<>();
final HashSet<String> 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<String> 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<String> container = new DefaultPulsarMessageListenerContainer<>(
pulsarConsumerFactory, pulsarContainerProperties);
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "dpmlct-013");
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
pulsarClient, prodConfig);
final PulsarTemplate<String> 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<String, Object> config = new HashMap<>();
final HashSet<String> 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<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
pulsarClient, config);
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
List<String> messages = new ArrayList<>();
pulsarContainerProperties.setMessageListener(
(PulsarRecordMessageListener<?>) (consumer, msg) -> messages.add((String) msg.getValue()));
pulsarContainerProperties.setSchema(Schema.STRING);
DefaultPulsarMessageListenerContainer<String> container = new DefaultPulsarMessageListenerContainer<>(
pulsarConsumerFactory, pulsarContainerProperties);
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "dpmlct-014");
final DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(
pulsarClient, prodConfig);
final PulsarTemplate<String> 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();
}
}