Consumer test cleanup

This commit is contained in:
Soby Chacko
2023-02-01 11:05:51 -05:00
parent 5dbb9a18c2
commit a7b0cf9d8b
3 changed files with 42 additions and 70 deletions

View File

@@ -29,7 +29,6 @@ import static org.mockito.Mockito.when;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicInteger;
@@ -57,9 +56,8 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
@Test
@SuppressWarnings("unchecked")
void happyPathErrorHandlingForRecordMessageListener() throws Exception {
Map<String, Object> config = new HashMap<>();
config.put("topicNames", Collections.singleton("default-error-handler-tests-1"));
config.put("subscriptionName", "default-error-handler-tests-sub-1");
Map<String, Object> config = Map.of("topicNames", Collections.singleton("default-error-handler-tests-1"),
"subscriptionName", "default-error-handler-tests-sub-1");
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
@@ -76,8 +74,7 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
pulsarContainerProperties.setMessageListener(messageListener);
pulsarContainerProperties.setSchema(Schema.STRING);
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "default-error-handler-tests-1");
Map<String, Object> prodConfig = Map.of("topicName", "default-error-handler-tests-1");
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
prodConfig);
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -108,9 +105,8 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
@Test
@SuppressWarnings("unchecked")
void errorHandlingForRecordMessageListenerWithTransientError() throws Exception {
Map<String, Object> config = new HashMap<>();
config.put("topicNames", Collections.singleton("default-error-handler-tests-2"));
config.put("subscriptionName", "default-error-handler-tests-sub-2");
Map<String, Object> config = Map.of("topicNames", Collections.singleton("default-error-handler-tests-2"),
"subscriptionName", "default-error-handler-tests-sub-2");
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
@@ -131,8 +127,7 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
pulsarContainerProperties.setMessageListener(messageListener);
pulsarContainerProperties.setSchema(Schema.STRING);
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "default-error-handler-tests-2");
Map<String, Object> prodConfig = Map.of("topicName", "default-error-handler-tests-2");
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
prodConfig);
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -157,9 +152,8 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
@Test
@SuppressWarnings("unchecked")
void everyOtherRecordThrowsNonTransientExceptionsRecordMessageListener() throws Exception {
Map<String, Object> config = new HashMap<>();
config.put("topicNames", Collections.singleton("default-error-handler-tests-3"));
config.put("subscriptionName", "default-error-handler-tests-sub-3");
Map<String, Object> config = Map.of("topicNames", Collections.singleton("default-error-handler-tests-3"),
"subscriptionName", "default-error-handler-tests-sub-3");
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
@@ -180,8 +174,7 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
pulsarContainerProperties.setMessageListener(messageListener);
pulsarContainerProperties.setSchema(Schema.INT32);
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "default-error-handler-tests-3");
Map<String, Object> prodConfig = Map.of("topicName", "default-error-handler-tests-3");
DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
prodConfig);
PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -216,9 +209,8 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
@Test
@SuppressWarnings("unchecked")
void batchRecordListenerFirstOneOnlyErrorAndRecover() throws Exception {
Map<String, Object> config = new HashMap<>();
config.put("topicNames", Collections.singleton("default-error-handler-tests-4"));
config.put("subscriptionName", "default-error-handler-tests-sub-4");
Map<String, Object> config = Map.of("topicNames", Collections.singleton("default-error-handler-tests-4"),
"subscriptionName", "default-error-handler-tests-sub-4");
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
@@ -260,8 +252,7 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
container.start();
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "default-error-handler-tests-4");
Map<String, Object> prodConfig = Map.of("topicName", "default-error-handler-tests-4");
DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
prodConfig);
PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -288,9 +279,8 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
@Test
@SuppressWarnings("unchecked")
void batchRecordListenerRecordFailsInTheMiddle() throws Exception {
Map<String, Object> config = new HashMap<>();
config.put("topicNames", Collections.singleton("default-error-handler-tests-5"));
config.put("subscriptionName", "default-error-handler-tests-sub-5");
Map<String, Object> config = Map.of("topicNames", Collections.singleton("default-error-handler-tests-5"),
"subscriptionName", "default-error-handler-tests-sub-5");
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
@@ -331,8 +321,7 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
container.start();
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "default-error-handler-tests-5");
Map<String, Object> prodConfig = Map.of("topicName", "default-error-handler-tests-5");
DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
prodConfig);
PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -358,9 +347,8 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
@Test
@SuppressWarnings("unchecked")
void batchRecordListenerRecordFailsTwiceInTheMiddle() throws Exception {
Map<String, Object> config = new HashMap<>();
config.put("topicNames", Collections.singleton("default-error-handler-tests-6"));
config.put("subscriptionName", "default-error-handler-tests-sub-6");
Map<String, Object> config = Map.of("topicNames", Collections.singleton("default-error-handler-tests-6"),
"subscriptionName", "default-error-handler-tests-sub-6");
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
@@ -401,8 +389,7 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
container.start();
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "default-error-handler-tests-6");
Map<String, Object> prodConfig = Map.of("topicName", "default-error-handler-tests-6");
DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
prodConfig);
PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -428,9 +415,8 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
@Test
@SuppressWarnings("unchecked")
void batchRecordListenerRecordFailsInTheMiddleButTransientError() throws Exception {
Map<String, Object> config = new HashMap<>();
config.put("topicNames", Collections.singleton("default-error-handler-tests-7"));
config.put("subscriptionName", "default-error-handler-tests-sub-7");
Map<String, Object> config = Map.of("topicNames", Collections.singleton("default-error-handler-tests-7"),
"subscriptionName", "default-error-handler-tests-sub-7");
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
@@ -477,8 +463,7 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
container.start();
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "default-error-handler-tests-7");
Map<String, Object> prodConfig = Map.of("topicName", "default-error-handler-tests-7");
DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
prodConfig);
PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -497,9 +482,8 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
@Test
@SuppressWarnings("unchecked")
void batchListenerFailsTransientErrorFollowedByNonTransient() throws Exception {
Map<String, Object> config = new HashMap<>();
config.put("topicNames", Collections.singleton("default-error-handler-tests-8"));
config.put("subscriptionName", "default-error-handler-tests-sub-8");
Map<String, Object> config = Map.of("topicNames", Collections.singleton("default-error-handler-tests-8"),
"subscriptionName", "default-error-handler-tests-sub-8");
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
@@ -549,8 +533,7 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain
container.start();
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "default-error-handler-tests-8");
Map<String, Object> prodConfig = Map.of("topicName", "default-error-handler-tests-8");
DefaultPulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
prodConfig);
PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);

View File

@@ -27,10 +27,8 @@ import static org.mockito.Mockito.verify;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.Condition;
@@ -65,8 +63,8 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
@Test
void basicDefaultConsumer() throws Exception {
Set<String> topics = Collections.singleton("dpmlct-012");
Map<String, Object> config = Map.of("topicNames", topics, "subscriptionName", "dpmlct-sb-012");
Map<String, Object> config = Map.of("topicNames", Collections.singleton("dpmlct-012"), "subscriptionName",
"dpmlct-sb-012");
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
@@ -81,8 +79,7 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
pulsarConsumerFactory, pulsarContainerProperties);
container.start();
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "dpmlct-012");
Map<String, Object> prodConfig = Map.of("topicName", "dpmlct-012");
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
prodConfig);
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -94,9 +91,8 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
@Test
void containerPauseAndResumeFeatureUsingWaitAndNotify() throws Exception {
Set<String> topics = Collections.singleton("containerPauseResumeWaitNotify-topic");
Map<String, Object> config = Map.of("topicNames", topics, "subscriptionName",
"containerPauseResumeWaitNotify-sub");
Map<String, Object> config = Map.of("topicNames", Collections.singleton("containerPauseResumeWaitNotify-topic"),
"subscriptionName", "containerPauseResumeWaitNotify-sub");
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
@@ -164,9 +160,8 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
@Test
void subscriptionInitialPositionEarliest() throws Exception {
Set<String> topics = Collections.singleton("dpmlct-013");
Map<String, Object> config = Map.of("topicNames", topics, "subscriptionName", "dpmlct-sb-013",
"subscriptionInitialPosition", SubscriptionInitialPosition.Earliest);
Map<String, Object> config = Map.of("topicNames", Collections.singleton("dpmlct-013"), "subscriptionName",
"dpmlct-sb-013", "subscriptionInitialPosition", SubscriptionInitialPosition.Earliest);
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
@@ -180,8 +175,7 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
DefaultPulsarMessageListenerContainer<String> container = new DefaultPulsarMessageListenerContainer<>(
pulsarConsumerFactory, pulsarContainerProperties);
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "dpmlct-013");
Map<String, Object> prodConfig = Map.of("topicName", "dpmlct-013");
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
prodConfig);
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -197,8 +191,8 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
@Test
void subscriptionInitialPositionDefaultLatest() throws Exception {
Set<String> topics = Collections.singleton("dpmlct-014");
Map<String, Object> config = Map.of("topicNames", topics, "subscriptionName", "dpmlct-sb-014");
Map<String, Object> config = Map.of("topicNames", Collections.singleton("dpmlct-014"), "subscriptionName",
"dpmlct-sb-014");
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
@@ -212,8 +206,7 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
DefaultPulsarMessageListenerContainer<String> container = new DefaultPulsarMessageListenerContainer<>(
pulsarConsumerFactory, pulsarContainerProperties);
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "dpmlct-014");
Map<String, Object> prodConfig = Map.of("topicName", "dpmlct-014");
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
prodConfig);
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -232,11 +225,10 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
@Test
void negativeAckRedeliveryBackoff() throws Exception {
Set<String> topics = Collections.singleton("dpmlct-015");
RedeliveryBackoff redeliveryBackoff = MultiplierRedeliveryBackoff.builder().minDelayMs(1000)
.maxDelayMs(5 * 1000).build();
Map<String, Object> config = Map.of("topicNames", topics, "subscriptionName", "dpmlct-sb-015",
"negativeAckRedeliveryBackoff", redeliveryBackoff);
Map<String, Object> config = Map.of("topicNames", Collections.singleton("dpmlct-015"), "subscriptionName",
"dpmlct-sb-015", "negativeAckRedeliveryBackoff", redeliveryBackoff);
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
@@ -280,11 +272,10 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
@Test
void deadLetterPolicyDefault() throws Exception {
Set<String> topics = Collections.singleton("dpmlct-016");
DeadLetterPolicy deadLetterPolicy = DeadLetterPolicy.builder().maxRedeliverCount(1)
.deadLetterTopic("dpmlct-016-dlq-topic").build();
Map<String, Object> config = Map.of("topicNames", topics, "subscriptionName", "dpmlct-sb-016",
"ackTimeoutMillis", 1, "deadLetterPolicy", deadLetterPolicy);
Map<String, Object> config = Map.of("topicNames", Collections.singleton("dpmlct-016"), "subscriptionName",
"dpmlct-sb-016", "ackTimeoutMillis", 1, "deadLetterPolicy", deadLetterPolicy);
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
@@ -336,11 +327,10 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS
@Test
void deadLetterPolicyCustom() throws Exception {
Set<String> topics = Collections.singleton("dpmlct-017");
DeadLetterPolicy deadLetterPolicy = DeadLetterPolicy.builder().maxRedeliverCount(5).deadLetterTopic("dlq-topic")
.build();
Map<String, Object> config = Map.of("topicNames", topics, "subscriptionName", "dpmlct-sb-016",
"ackTimeoutMillis", 1, "deadLetterPolicy", deadLetterPolicy);
Map<String, Object> config = Map.of("topicNames", Collections.singleton("dpmlct-017"), "subscriptionName",
"dpmlct-sb-016", "ackTimeoutMillis", 1, "deadLetterPolicy", deadLetterPolicy);
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();

View File

@@ -103,8 +103,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport {
@Bean
public PulsarProducerFactory<String> pulsarProducerFactory(PulsarClient pulsarClient) {
Map<String, Object> config = new HashMap<>();
config.put("topicName", "foo-1");
Map<String, Object> config = Map.of("topicName", "foo-1");
return new DefaultPulsarProducerFactory<>(pulsarClient, config);
}