From 27bba9816e42127c78378cf4faef951f3d671a26 Mon Sep 17 00:00:00 2001 From: Christophe Bornet Date: Thu, 17 Nov 2022 21:21:09 +0100 Subject: [PATCH] Remove final keyword in methods (#213) --- ...arListenerAnnotationBeanPostProcessor.java | 2 +- ...arListenerAnnotationBeanPostProcessor.java | 2 +- .../config/MethodPulsarListenerEndpoint.java | 20 +-- .../MethodReactivePulsarListenerEndpoint.java | 21 ++- .../core/CachingPulsarProducerFactory.java | 2 +- .../pulsar/core/PulsarTemplate.java | 4 +- .../DefaultReactivePulsarSenderFactory.java | 4 +- .../core/reactive/ReactivePulsarTemplate.java | 4 +- .../DefaultPulsarConsumerErrorHandler.java | 4 +- ...DefaultPulsarMessageListenerContainer.java | 35 +++-- ...rBatchMessagingMessageListenerAdapter.java | 11 +- .../PulsarMessagingMessageConverter.java | 2 +- .../core/ConsumerAcknowledgmentTests.java | 120 +++++++-------- .../pulsar/core/FailoverConsumerTests.java | 18 +-- ...efaultPulsarConsumerErrorHandlerTests.java | 138 +++++++++--------- ...ltPulsarMessageListenerContainerTests.java | 90 ++++++------ .../pulsar/listener/PulsarListenerTests.java | 26 ++-- .../reactive/ReactivePulsarListenerTests.java | 6 +- 18 files changed, 253 insertions(+), 256 deletions(-) diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarListenerAnnotationBeanPostProcessor.java b/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarListenerAnnotationBeanPostProcessor.java index cfaf560c..f7f543ac 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarListenerAnnotationBeanPostProcessor.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarListenerAnnotationBeanPostProcessor.java @@ -244,7 +244,7 @@ public class PulsarListenerAnnotationBeanPostProcessor } @Override - public Object postProcessAfterInitialization(final Object bean, final String beanName) throws BeansException { + public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException { if (!this.nonAnnotatedClasses.contains(bean.getClass())) { Class targetClass = AopUtils.getTargetClass(bean); Map> annotatedMethods = MethodIntrospector.selectMethods(targetClass, diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/ReactivePulsarListenerAnnotationBeanPostProcessor.java b/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/ReactivePulsarListenerAnnotationBeanPostProcessor.java index 7105bc87..4f28a903 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/ReactivePulsarListenerAnnotationBeanPostProcessor.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/ReactivePulsarListenerAnnotationBeanPostProcessor.java @@ -242,7 +242,7 @@ public class ReactivePulsarListenerAnnotationBeanPostProcessor } @Override - public Object postProcessAfterInitialization(final Object bean, final String beanName) throws BeansException { + public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException { if (!this.nonAnnotatedClasses.contains(bean.getClass())) { Class targetClass = AopUtils.getTargetClass(bean); Map> annotatedMethods = MethodIntrospector.selectMethods(targetClass, diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarListenerEndpoint.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarListenerEndpoint.java index bf17e6ba..77f3e9b0 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarListenerEndpoint.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarListenerEndpoint.java @@ -120,7 +120,7 @@ public class MethodPulsarListenerEndpoint extends AbstractPulsarListenerEndpo Assert.state(this.messageHandlerMethodFactory != null, "Could not create message listener - MessageHandlerMethodFactory not set"); PulsarMessagingMessageListenerAdapter messageListener = createMessageListenerInstance(messageConverter); - final HandlerAdapter handlerMethod = configureListenerAdapter(messageListener); + HandlerAdapter handlerMethod = configureListenerAdapter(messageListener); messageListener.setHandlerMethod(handlerMethod); // Since we have access to the handler method here, check if we can type infer the @@ -128,14 +128,14 @@ public class MethodPulsarListenerEndpoint extends AbstractPulsarListenerEndpo // TODO: filter out the payload type by excluding Consumer, Message, Messages etc. - final MethodParameter[] methodParameters = handlerMethod.getInvokerHandlerMethod().getMethodParameters(); + MethodParameter[] methodParameters = handlerMethod.getInvokerHandlerMethod().getMethodParameters(); MethodParameter messageParameter = null; - final Optional parameter = Arrays.stream(methodParameters) + Optional parameter = Arrays.stream(methodParameters) .filter(methodParameter1 -> !methodParameter1.getParameterType().equals(Consumer.class) || !methodParameter1.getParameterType().equals(Acknowledgement.class) || !methodParameter1.hasParameterAnnotation(Header.class)) .findFirst(); - final long count = Arrays.stream(methodParameters) + long count = Arrays.stream(methodParameters) .filter(methodParameter1 -> !methodParameter1.getParameterType().equals(Consumer.class) && !methodParameter1.getParameterType().equals(Acknowledgement.class) && !methodParameter1.hasParameterAnnotation(Header.class)) @@ -145,9 +145,9 @@ public class MethodPulsarListenerEndpoint extends AbstractPulsarListenerEndpo messageParameter = parameter.get(); } - final ConcurrentPulsarMessageListenerContainer containerInstance = (ConcurrentPulsarMessageListenerContainer) container; - final PulsarContainerProperties pulsarContainerProperties = containerInstance.getContainerProperties(); - final SchemaType schemaType = pulsarContainerProperties.getSchemaType(); + ConcurrentPulsarMessageListenerContainer containerInstance = (ConcurrentPulsarMessageListenerContainer) container; + PulsarContainerProperties pulsarContainerProperties = containerInstance.getContainerProperties(); + SchemaType schemaType = pulsarContainerProperties.getSchemaType(); if (schemaType != SchemaType.NONE) { switch (schemaType) { case STRING -> pulsarContainerProperties.setSchema(Schema.STRING); @@ -193,7 +193,7 @@ public class MethodPulsarListenerEndpoint extends AbstractPulsarListenerEndpo } } } - final SchemaType type = pulsarContainerProperties.getSchema().getSchemaInfo().getType(); + SchemaType type = pulsarContainerProperties.getSchema().getSchemaInfo().getType(); pulsarContainerProperties.setSchemaType(type); container.setNegativeAckRedeliveryBackoff(this.negativeAckRedeliveryBackoff); @@ -206,7 +206,7 @@ public class MethodPulsarListenerEndpoint extends AbstractPulsarListenerEndpo private Schema getMessageSchema(MethodParameter messageParameter, Function, Schema> schemaFactory) { ResolvableType messageType = resolvableType(messageParameter); - final Class messageClass = messageType.getRawClass(); + Class messageClass = messageType.getRawClass(); return schemaFactory.apply(messageClass); } @@ -221,7 +221,7 @@ public class MethodPulsarListenerEndpoint extends AbstractPulsarListenerEndpo private ResolvableType resolvableType(MethodParameter methodParameter) { ResolvableType resolvableType = ResolvableType.forMethodParameter(methodParameter); - final Class rawClass = resolvableType.getRawClass(); + Class rawClass = resolvableType.getRawClass(); if (rawClass != null && isContainerType(rawClass)) { resolvableType = resolvableType.getGeneric(0); } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/reactive/MethodReactivePulsarListenerEndpoint.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/reactive/MethodReactivePulsarListenerEndpoint.java index 8d50fd06..ae62fe2a 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/reactive/MethodReactivePulsarListenerEndpoint.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/reactive/MethodReactivePulsarListenerEndpoint.java @@ -114,7 +114,7 @@ public class MethodReactivePulsarListenerEndpoint extends AbstractReactivePul Assert.state(this.messageHandlerMethodFactory != null, "Could not create message listener - MessageHandlerMethodFactory not set"); PulsarMessagingMessageListenerAdapter messageListener = createMessageListenerInstance(messageConverter); - final HandlerAdapter handlerMethod = configureListenerAdapter(messageListener); + HandlerAdapter handlerMethod = configureListenerAdapter(messageListener); messageListener.setHandlerMethod(handlerMethod); // Since we have access to the handler method here, check if we can type infer the @@ -122,14 +122,14 @@ public class MethodReactivePulsarListenerEndpoint extends AbstractReactivePul // TODO: filter out the payload type by excluding Consumer, Message, Messages etc. - final MethodParameter[] methodParameters = handlerMethod.getInvokerHandlerMethod().getMethodParameters(); + MethodParameter[] methodParameters = handlerMethod.getInvokerHandlerMethod().getMethodParameters(); MethodParameter messageParameter = null; - final Optional parameter = Arrays.stream(methodParameters) + Optional parameter = Arrays.stream(methodParameters) .filter(methodParameter1 -> !methodParameter1.getParameterType().equals(Consumer.class) || !methodParameter1.getParameterType().equals(Acknowledgement.class) || !methodParameter1.hasParameterAnnotation(Header.class)) .findFirst(); - final long count = Arrays.stream(methodParameters) + long count = Arrays.stream(methodParameters) .filter(methodParameter1 -> !methodParameter1.getParameterType().equals(Consumer.class) && !methodParameter1.getParameterType().equals(Acknowledgement.class) && !methodParameter1.hasParameterAnnotation(Header.class)) @@ -139,10 +139,9 @@ public class MethodReactivePulsarListenerEndpoint extends AbstractReactivePul messageParameter = parameter.get(); } - final DefaultReactivePulsarMessageListenerContainer containerInstance = (DefaultReactivePulsarMessageListenerContainer) container; - final ReactivePulsarContainerProperties pulsarContainerProperties = containerInstance - .getContainerProperties(); - final SchemaType schemaType = pulsarContainerProperties.getSchemaType(); + DefaultReactivePulsarMessageListenerContainer containerInstance = (DefaultReactivePulsarMessageListenerContainer) container; + ReactivePulsarContainerProperties pulsarContainerProperties = containerInstance.getContainerProperties(); + SchemaType schemaType = pulsarContainerProperties.getSchemaType(); if (schemaType != SchemaType.NONE) { switch (schemaType) { case STRING -> pulsarContainerProperties.setSchema((Schema) Schema.STRING); @@ -188,7 +187,7 @@ public class MethodReactivePulsarListenerEndpoint extends AbstractReactivePul } } } - final SchemaType type = pulsarContainerProperties.getSchema().getSchemaInfo().getType(); + SchemaType type = pulsarContainerProperties.getSchema().getSchemaInfo().getType(); pulsarContainerProperties.setSchemaType(type); ReactiveMessageConsumerBuilderCustomizer customizer1 = b -> b.deadLetterPolicy(this.deadLetterPolicy); @@ -204,7 +203,7 @@ public class MethodReactivePulsarListenerEndpoint extends AbstractReactivePul private Schema getMessageSchema(MethodParameter messageParameter, Function, Schema> schemaFactory) { ResolvableType messageType = resolvableType(messageParameter); - final Class messageClass = messageType.getRawClass(); + Class messageClass = messageType.getRawClass(); return schemaFactory.apply(messageClass); } @@ -219,7 +218,7 @@ public class MethodReactivePulsarListenerEndpoint extends AbstractReactivePul private ResolvableType resolvableType(MethodParameter methodParameter) { ResolvableType resolvableType = ResolvableType.forMethodParameter(methodParameter); - final Class rawClass = resolvableType.getRawClass(); + Class rawClass = resolvableType.getRawClass(); if (rawClass != null && isContainerType(rawClass)) { resolvableType = resolvableType.getGeneric(0); } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/CachingPulsarProducerFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/CachingPulsarProducerFactory.java index a0b9d1d2..f91c94f9 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/CachingPulsarProducerFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/CachingPulsarProducerFactory.java @@ -95,7 +95,7 @@ public class CachingPulsarProducerFactory extends DefaultPulsarProducerFactor @Override protected Producer doCreateProducer(Schema schema, @Nullable String topic, @Nullable Collection encryptionKeys, @Nullable List> customizers) { - final String topicName = ProducerUtils.resolveTopicName(topic, this); + String topicName = ProducerUtils.resolveTopicName(topic, this); ProducerCacheKey producerCacheKey = new ProducerCacheKey<>(schema, topicName, encryptionKeys == null ? null : new HashSet<>(encryptionKeys), customizers); return this.producerCache.get(producerCacheKey, diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarTemplate.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarTemplate.java index ebf9b0d2..e092914c 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarTemplate.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarTemplate.java @@ -158,14 +158,14 @@ public class PulsarTemplate implements PulsarOperations, BeanNameAware { @Nullable Collection encryptionKeys, T message, @Nullable TypedMessageBuilderCustomizer typedMessageBuilderCustomizer, @Nullable ProducerBuilderCustomizer producerCustomizer) throws PulsarClientException { - final String topicName = ProducerUtils.resolveTopicName(topic, this.producerFactory); + String topicName = ProducerUtils.resolveTopicName(topic, this.producerFactory); this.logger.trace(() -> String.format("Sending msg to '%s' topic", topicName)); PulsarMessageSenderContext senderContext = PulsarMessageSenderContext.newContext(topicName, this.beanName); Observation observation = newObservation(senderContext); try { observation.start(); - final Producer producer = prepareProducerForSend(topic, message, encryptionKeys, producerCustomizer); + Producer producer = prepareProducerForSend(topic, message, encryptionKeys, producerCustomizer); TypedMessageBuilder messageBuilder = producer.newMessage().value(message); if (typedMessageBuilderCustomizer != null) { typedMessageBuilderCustomizer.customize(messageBuilder); diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/DefaultReactivePulsarSenderFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/DefaultReactivePulsarSenderFactory.java index e9117bf4..9004221f 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/DefaultReactivePulsarSenderFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/DefaultReactivePulsarSenderFactory.java @@ -77,9 +77,9 @@ public class DefaultReactivePulsarSenderFactory implements ReactivePulsarSend private ReactiveMessageSender doCreateReactiveMessageSender(String topic, Schema schema, List> customizers) { - final String resolvedTopic = ReactiveMessageSenderUtils.resolveTopicName(topic, this); + String resolvedTopic = ReactiveMessageSenderUtils.resolveTopicName(topic, this); this.logger.trace(() -> String.format("Creating reactive message sender for '%s' topic", resolvedTopic)); - final ReactiveMessageSenderBuilder sender = this.reactivePulsarClient.messageSender(schema); + ReactiveMessageSenderBuilder sender = this.reactivePulsarClient.messageSender(schema); sender.applySpec(this.reactiveMessageSenderSpec); sender.topic(resolvedTopic); if (this.reactiveMessageSenderCache != null) { diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplate.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplate.java index b1dba2b2..da622040 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplate.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplate.java @@ -90,7 +90,7 @@ public class ReactivePulsarTemplate implements ReactivePulsarOperations { private Mono doSend(String topic, T message, MessageSpecBuilderCustomizer messageSpecBuilderCustomizer, ReactiveMessageSenderBuilderCustomizer customizer) { - final String topicName = ReactiveMessageSenderUtils.resolveTopicName(topic, this.reactiveMessageSenderFactory); + String topicName = ReactiveMessageSenderUtils.resolveTopicName(topic, this.reactiveMessageSenderFactory); this.logger.trace(() -> String.format("Sending reactive message to '%s' topic", topicName)); ReactiveMessageSender sender = createMessageSender(topic, message, customizer); return sender.sendOne(getMessageSpec(messageSpecBuilderCustomizer, message)).doOnError( @@ -100,7 +100,7 @@ public class ReactivePulsarTemplate implements ReactivePulsarOperations { } private Flux doSendMany(String topic, Flux messages) { - final String topicName = ReactiveMessageSenderUtils.resolveTopicName(topic, this.reactiveMessageSenderFactory); + String topicName = ReactiveMessageSenderUtils.resolveTopicName(topic, this.reactiveMessageSenderFactory); this.logger.trace(() -> String.format("Sending reactive messages to '%s' topic", topicName)); if (this.schema != null) { diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarConsumerErrorHandler.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarConsumerErrorHandler.java index f29ce419..b398cd06 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarConsumerErrorHandler.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarConsumerErrorHandler.java @@ -51,7 +51,7 @@ public class DefaultPulsarConsumerErrorHandler implements PulsarConsumerError @Override public boolean shouldRetryMessage(Exception exception, Message message) { - final Pair pair = this.backOffExecutionThreadLocal.get(); + Pair pair = this.backOffExecutionThreadLocal.get(); long nextBackOff; BackOffExecution backOffExecution; if (pair != null && pair.message.equals(message)) { @@ -85,7 +85,7 @@ public class DefaultPulsarConsumerErrorHandler implements PulsarConsumerError @SuppressWarnings("unchecked") public Message currentMessage() { // there is only one message tracked at any time. - final Pair pair = this.backOffExecutionThreadLocal.get(); + Pair pair = this.backOffExecutionThreadLocal.get(); if (pair == null) { return null; } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java index 0f430c9e..7df2cbcb 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java @@ -274,28 +274,28 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess private void populateAllNecessaryPropertiesIfNeedBe(Map currentProperties) { if (currentProperties.containsKey("topicNames")) { - final String topicsFromMap = (String) currentProperties.get("topicNames"); - final String[] topicNames = StringUtils.delimitedListToStringArray(topicsFromMap, ","); - final Set propertiesDefinedTopics = Set.of(topicNames); + String topicsFromMap = (String) currentProperties.get("topicNames"); + String[] topicNames = StringUtils.delimitedListToStringArray(topicsFromMap, ","); + Set propertiesDefinedTopics = Set.of(topicNames); if (!propertiesDefinedTopics.isEmpty()) { currentProperties.put("topicNames", propertiesDefinedTopics); } } if (!currentProperties.containsKey("subscriptionType")) { - final SubscriptionType subscriptionType = this.containerProperties.getSubscriptionType(); + SubscriptionType subscriptionType = this.containerProperties.getSubscriptionType(); if (subscriptionType != null) { currentProperties.put("subscriptionType", subscriptionType); } } if (!currentProperties.containsKey("topicNames")) { - final String[] topics = this.containerProperties.getTopics(); - final Set listenerDefinedTopics = new HashSet<>(Arrays.stream(topics).toList()); + String[] topics = this.containerProperties.getTopics(); + Set listenerDefinedTopics = new HashSet<>(Arrays.stream(topics).toList()); if (!listenerDefinedTopics.isEmpty()) { currentProperties.put("topicNames", listenerDefinedTopics); } } if (!currentProperties.containsKey("topicsPattern")) { - final String topicsPattern = this.containerProperties.getTopicsPattern(); + String topicsPattern = this.containerProperties.getTopicsPattern(); if (topicsPattern != null) { currentProperties.put("topicsPattern", topicsPattern); } @@ -305,15 +305,15 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess currentProperties.put("subscriptionName", this.containerProperties.getSubscriptionName()); } } - final RedeliveryBackoff negativeAckRedeliveryBackoff = DefaultPulsarMessageListenerContainer.this.negativeAckRedeliveryBackoff; + RedeliveryBackoff negativeAckRedeliveryBackoff = DefaultPulsarMessageListenerContainer.this.negativeAckRedeliveryBackoff; if (negativeAckRedeliveryBackoff != null) { currentProperties.put("negativeAckRedeliveryBackoff", negativeAckRedeliveryBackoff); } - final RedeliveryBackoff ackTimeoutRedeliveryBackoff = DefaultPulsarMessageListenerContainer.this.ackTimeoutRedeliveryBackoff; + RedeliveryBackoff ackTimeoutRedeliveryBackoff = DefaultPulsarMessageListenerContainer.this.ackTimeoutRedeliveryBackoff; if (ackTimeoutRedeliveryBackoff != null) { currentProperties.put("ackTimeoutRedeliveryBackoff", ackTimeoutRedeliveryBackoff); } - final DeadLetterPolicy deadLetterPolicy = DefaultPulsarMessageListenerContainer.this.deadLetterPolicy; + DeadLetterPolicy deadLetterPolicy = DefaultPulsarMessageListenerContainer.this.deadLetterPolicy; if (deadLetterPolicy != null) { currentProperties.put("deadLetterPolicy", deadLetterPolicy); } @@ -380,8 +380,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess this.consumer.acknowledge(messages); } else { - final Stream> stream = StreamSupport.stream(messages.spliterator(), - true); + Stream> stream = StreamSupport.stream(messages.spliterator(), true); Message last = stream.reduce((a, b) -> b).orElse(null); this.consumer.acknowledgeCumulative(last); } @@ -493,7 +492,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess PulsarBatchListenerFailedException pulsarBatchListenerFailedException = (PulsarBatchListenerFailedException) exception; Message pulsarMessage = getPulsarMessageCausedTheException(pulsarBatchListenerFailedException); - final Message theCurrentPulsarMessageTracked = this.pulsarConsumerErrorHandler.currentMessage(); + Message theCurrentPulsarMessageTracked = this.pulsarConsumerErrorHandler.currentMessage(); // Previous message in error handled during retry but another msg in sublist // caused error; // resetting state in order to track it @@ -507,10 +506,10 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess // handled on the retry. Otherwise, if we are out of retries then the sublist // does not include // the message in error (it instead gets recovered). - final int indexOfFailedMessage = messageList.indexOf(pulsarMessage); + int indexOfFailedMessage = messageList.indexOf(pulsarMessage); messageList = messageList.subList(indexOfFailedMessage, messageList.size()); - final boolean toBeRetried = this.pulsarConsumerErrorHandler - .shouldRetryMessage(pulsarBatchListenerFailedException, pulsarMessage); + boolean toBeRetried = this.pulsarConsumerErrorHandler.shouldRetryMessage(pulsarBatchListenerFailedException, + pulsarMessage); if (toBeRetried) { inRetryMode.set(true); } @@ -535,7 +534,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess } private void invokeRecordListenerErrorHandler(AtomicBoolean inRetryMode, Message message, Exception e) { - final boolean toBeRetried = this.pulsarConsumerErrorHandler.shouldRetryMessage(e, message); + boolean toBeRetried = this.pulsarConsumerErrorHandler.shouldRetryMessage(e, message); if (toBeRetried) { inRetryMode.set(true); } @@ -576,7 +575,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess this.consumer.acknowledge(messages); } else { - final Stream> stream = StreamSupport.stream(messages.spliterator(), true); + Stream> stream = StreamSupport.stream(messages.spliterator(), true); Message last = stream.reduce((a, b) -> b).orElse(null); this.consumer.acknowledgeCumulative(last); } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/adapter/PulsarBatchMessagingMessageListenerAdapter.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/adapter/PulsarBatchMessagingMessageListenerAdapter.java index 87b1eb35..e00bc18f 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/adapter/PulsarBatchMessagingMessageListenerAdapter.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/adapter/PulsarBatchMessagingMessageListenerAdapter.java @@ -79,7 +79,7 @@ public class PulsarBatchMessagingMessageListenerAdapter extends PulsarMessagi else if (isPulsarMessageList() && isHeaderFound()) { // List, // @Header List> messages = toSpringMessages(consumer, msg); - final Map> aggregatedHeaders = withAggregatedHeaders(messages); + Map> aggregatedHeaders = withAggregatedHeaders(messages); List list1 = new ArrayList<>(msg); message = MessageBuilder.withPayload(list1).copyHeaders(aggregatedHeaders).build(); } @@ -89,7 +89,7 @@ public class PulsarBatchMessagingMessageListenerAdapter extends PulsarMessagi } else if (isMessageList() && isHeaderFound()) { // List, @Header List> messages = toSpringMessages(consumer, msg); - final Map> aggregatedHeaders = withAggregatedHeaders(messages); + Map> aggregatedHeaders = withAggregatedHeaders(messages); message = MessageBuilder.withPayload(messages).copyHeaders(aggregatedHeaders).build(); } else if (this.isSimpleExtraction()) { // List @@ -99,7 +99,7 @@ public class PulsarBatchMessagingMessageListenerAdapter extends PulsarMessagi } else if (isHeaderFound()) { // List, @Header List> messages = toSpringMessages(consumer, msg); - final Map> aggregatedHeaders = withAggregatedHeaders(messages); + Map> aggregatedHeaders = withAggregatedHeaders(messages); List list = new ArrayList<>(msg.size()); msg.stream().iterator().forEachRemaining(vMessage -> list.add(vMessage.getValue())); message = MessageBuilder.withPayload(list).copyHeaders(aggregatedHeaders).build(); @@ -123,7 +123,7 @@ public class PulsarBatchMessagingMessageListenerAdapter extends PulsarMessagi } private Map> withAggregatedHeaders(List> messages) { - final Map> aggregatedHeaders = new HashMap<>(); + Map> aggregatedHeaders = new HashMap<>(); for (Message message : messages) { message.getHeaders().forEach((s, o) -> { List objects = aggregatedHeaders.computeIfAbsent(s, k -> new ArrayList<>()); @@ -139,8 +139,7 @@ public class PulsarBatchMessagingMessageListenerAdapter extends PulsarMessagi return messages; } - protected void invoke(Object records, Consumer consumer, final Message message, - Acknowledgement acknowledgement) { + protected void invoke(Object records, Consumer consumer, Message message, Acknowledgement acknowledgement) { try { invokeHandler(records, message, consumer, acknowledgement); } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/support/converter/PulsarMessagingMessageConverter.java b/spring-pulsar/src/main/java/org/springframework/pulsar/support/converter/PulsarMessagingMessageConverter.java index 4e2386e5..0a12c1e7 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/support/converter/PulsarMessagingMessageConverter.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/support/converter/PulsarMessagingMessageConverter.java @@ -49,7 +49,7 @@ public class PulsarMessagingMessageConverter implements PulsarRecordMessageCo @Override public Message toMessage(org.apache.pulsar.client.api.Message record, Consumer consumer, Type type) { - final Map messageHeaders = new HashMap<>(); + Map messageHeaders = new HashMap<>(); this.pulsarMessageHeaderMapper.toHeaders(record, messageHeaders); Message message = MessageBuilder.createMessage(extractAndConvertValue(record), new MessageHeaders(messageHeaders)); 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 284a4f70..bf714402 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 @@ -64,13 +64,13 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { @Test void testRecordAck() throws Exception { Map config = new HashMap<>(); - final Set strings = new HashSet<>(); + Set strings = new HashSet<>(); strings.add("cons-ack-tests-011"); config.put("topicNames", strings); config.put("subscriptionName", "cons-ack-tests-sb-011"); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = spy( + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = spy( new DefaultPulsarConsumerFactory<>(pulsarClient, config)); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); @@ -80,7 +80,7 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { pulsarContainerProperties.setAckMode(AckMode.RECORD); DefaultPulsarMessageListenerContainer container = new DefaultPulsarMessageListenerContainer<>( pulsarConsumerFactory, pulsarContainerProperties); - final Consumer containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container); + Consumer containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container); CountDownLatch latch = new CountDownLatch(10); @@ -91,9 +91,9 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { Map prodConfig = new HashMap<>(); prodConfig.put("topicName", "cons-ack-tests-011"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); for (int i = 0; i < 10; i++) { pulsarTemplate.sendAsync("hello john doe"); } @@ -105,13 +105,13 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { @Test void testBatchAck() throws Exception { Map config = new HashMap<>(); - final Set strings = new HashSet<>(); + Set strings = new HashSet<>(); strings.add("cons-ack-tests-012"); config.put("topicNames", strings); config.put("subscriptionName", "cons-ack-tests-sb-012"); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = spy( + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = spy( new DefaultPulsarConsumerFactory<>(pulsarClient, config)); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); @@ -121,13 +121,13 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { pulsarContainerProperties.setSchema(Schema.STRING); DefaultPulsarMessageListenerContainer container = new DefaultPulsarMessageListenerContainer<>( pulsarConsumerFactory, pulsarContainerProperties); - final Consumer containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container); + Consumer containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container); Map prodConfig = new HashMap<>(); prodConfig.put("topicName", "cons-ack-tests-012"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); for (int i = 0; i < 10; i++) { pulsarTemplate.sendAsync("hello john doe"); } @@ -145,9 +145,9 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { Map config = new HashMap<>(); config.put("topicNames", Collections.singleton("cons-ack-tests-013")); config.put("subscriptionName", "cons-ack-tests-sb-013"); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = spy( + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = spy( new DefaultPulsarConsumerFactory<>(pulsarClient, config)); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); @@ -162,7 +162,7 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { pulsarContainerProperties.setSchema(Schema.STRING); DefaultPulsarMessageListenerContainer container = new DefaultPulsarMessageListenerContainer<>( pulsarConsumerFactory, pulsarContainerProperties); - final Consumer containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container); + Consumer containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container); AtomicInteger ackCallCount = new AtomicInteger(0); doAnswer(invocation -> { @@ -172,9 +172,9 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { Map prodConfig = new HashMap<>(); prodConfig.put("topicName", "cons-ack-tests-013"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); for (int i = 0; i < 10; i++) { pulsarTemplate.sendAsync("hello john doe"); } @@ -186,7 +186,7 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { await().atMost(Duration.ofSeconds(10)) .untilAsserted(() -> verify(containerConsumer, times(5)).negativeAcknowledge(any(Message.class))); - final int ackCalls = ackCallCount.get(); + int ackCalls = ackCallCount.get(); if (ackCalls < 5) { await().atMost(Duration.ofSeconds(10)) .untilAsserted(() -> verify(containerConsumer, atMost(4)).acknowledge(any(MessageId.class))); @@ -209,17 +209,17 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { @SuppressWarnings("unchecked") void testManualAckForRecordListener() throws Exception { Map config = new HashMap<>(); - final Set strings = new HashSet<>(); + Set strings = new HashSet<>(); strings.add("cons-ack-tests-014"); config.put("topicNames", strings); config.put("subscriptionName", "cons-ack-tests-sb-014"); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = spy( + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = spy( new DefaultPulsarConsumerFactory<>(pulsarClient, config)); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); - final List acksObjects = new ArrayList<>(); + List acksObjects = new ArrayList<>(); PulsarAcknowledgingMessageListener pulsarAcknowledgingMessageListener = (consumer, msg, acknowledgement) -> { acksObjects.add(acknowledgement); acknowledgement.acknowledge(); @@ -231,7 +231,7 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { pulsarContainerProperties.setAckMode(AckMode.MANUAL); DefaultPulsarMessageListenerContainer container = new DefaultPulsarMessageListenerContainer<>( pulsarConsumerFactory, pulsarContainerProperties); - final Consumer containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container); + Consumer containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container); CountDownLatch latch = new CountDownLatch(10); @@ -242,9 +242,9 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { Map prodConfig = new HashMap<>(); prodConfig.put("topicName", "cons-ack-tests-014"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); for (int i = 0; i < 10; i++) { pulsarTemplate.sendAsync("hello john doe"); } @@ -263,13 +263,13 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { @SuppressWarnings("unchecked") void testBatchAckForBatchListener() throws Exception { Map config = new HashMap<>(); - final Set strings = new HashSet<>(); + Set strings = new HashSet<>(); strings.add("cons-ack-tests-015"); config.put("topicNames", strings); config.put("subscriptionName", "cons-ack-tests-sb-015"); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = spy( + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = spy( new DefaultPulsarConsumerFactory<>(pulsarClient, config)); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); @@ -277,7 +277,7 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { pulsarContainerProperties.setBatchTimeoutMillis(60_000); pulsarContainerProperties.setBatchListener(true); CountDownLatch latch = new CountDownLatch(1); - final PulsarBatchMessageListener pulsarBatchMessageListener = mock(PulsarBatchMessageListener.class); + PulsarBatchMessageListener pulsarBatchMessageListener = mock(PulsarBatchMessageListener.class); doAnswer(invocation -> { latch.countDown(); @@ -288,13 +288,13 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { pulsarContainerProperties.setSchema(Schema.STRING); DefaultPulsarMessageListenerContainer container = new DefaultPulsarMessageListenerContainer<>( pulsarConsumerFactory, pulsarContainerProperties); - final Consumer containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container); + Consumer containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container); Map prodConfig = new HashMap<>(); prodConfig.put("topicName", "cons-ack-tests-015"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); for (int i = 0; i < 10; i++) { pulsarTemplate.sendAsync("hello john doe"); } @@ -311,20 +311,20 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { @SuppressWarnings("unchecked") void testBatchNackForEntireBatchWhenUsingBatchListener() throws Exception { Map config = new HashMap<>(); - final Set strings = new HashSet<>(); + Set strings = new HashSet<>(); strings.add("cons-ack-tests-016"); config.put("topicNames", strings); config.put("subscriptionName", "cons-ack-tests-sb-016"); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = spy( + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = spy( new DefaultPulsarConsumerFactory<>(pulsarClient, config)); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties.setMaxNumMessages(10); pulsarContainerProperties.setBatchTimeoutMillis(60_000); pulsarContainerProperties.setBatchListener(true); - final PulsarBatchMessageListener pulsarBatchMessageListener = mock(PulsarBatchMessageListener.class); + PulsarBatchMessageListener pulsarBatchMessageListener = mock(PulsarBatchMessageListener.class); CountDownLatch latch = new CountDownLatch(1); doAnswer(invocation -> { @@ -336,13 +336,13 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { pulsarContainerProperties.setSchema(Schema.STRING); DefaultPulsarMessageListenerContainer container = new DefaultPulsarMessageListenerContainer<>( pulsarConsumerFactory, pulsarContainerProperties); - final Consumer containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container); + Consumer containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container); Map prodConfig = new HashMap<>(); prodConfig.put("topicName", "cons-ack-tests-016"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); for (int i = 0; i < 10; i++) { pulsarTemplate.sendAsync("hello john doe"); } @@ -361,13 +361,13 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { Map config = new HashMap<>(); config.put("topicNames", Set.of("duplicate-message-test")); config.put("subscriptionName", "duplicate-sub-1"); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( - pulsarClient, config); + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, + config); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); - final AtomicInteger counter1 = new AtomicInteger(0); + AtomicInteger counter1 = new AtomicInteger(0); pulsarContainerProperties.setMessageListener((PulsarRecordMessageListener) (consumer, msg) -> { counter1.getAndIncrement(); }); @@ -377,9 +377,9 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { container1.start(); Map prodConfig = Collections.singletonMap("topicName", "duplicate-message-test"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); pulsarTemplate.send("hello john doe"); while (counter1.get() == 0) { @@ -391,7 +391,7 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport { // unacked message. container1.stop(); - final AtomicInteger counter2 = new AtomicInteger(0); + AtomicInteger counter2 = new AtomicInteger(0); pulsarContainerProperties.setMessageListener((PulsarRecordMessageListener) (consumer, msg) -> { counter2.getAndIncrement(); }); diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/FailoverConsumerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/FailoverConsumerTests.java index ea37e7b4..abf1491f 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/FailoverConsumerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/FailoverConsumerTests.java @@ -54,14 +54,14 @@ class FailoverConsumerTests implements PulsarTestContainerSupport { admin.topics().createPartitionedTopic(topicName, numPartitions); Map config = new HashMap<>(); - final HashSet topics = new HashSet<>(); + HashSet topics = new HashSet<>(); topics.add("my-part-topic-1"); config.put("topicNames", topics); config.put("subscriptionName", "my-part-subscription-1"); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( - pulsarClient, config); + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, + config); CountDownLatch latch = new CountDownLatch(3); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties @@ -80,9 +80,9 @@ class FailoverConsumerTests implements PulsarTestContainerSupport { Map prodConfig = new HashMap<>(); prodConfig.put("topicName", "my-part-topic-1"); prodConfig.put("messageRoutingMode", MessageRoutingMode.CustomPartition); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); pulsarTemplate.newMessage("hello john doe") .withProducerCustomizer(builder -> builder.messageRouter(new FooRouter())).sendAsync(); @@ -91,7 +91,7 @@ class FailoverConsumerTests implements PulsarTestContainerSupport { pulsarTemplate.newMessage("hello buzz doe") .withProducerCustomizer(builder -> builder.messageRouter(new BuzzRouter())).sendAsync(); - final boolean await = latch.await(10, TimeUnit.SECONDS); + boolean await = latch.await(10, TimeUnit.SECONDS); assertThat(await).isTrue(); } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarConsumerErrorHandlerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarConsumerErrorHandlerTests.java index c2d6da89..bb8ce07c 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarConsumerErrorHandlerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarConsumerErrorHandlerTests.java @@ -61,10 +61,10 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain config.put("topicNames", Collections.singleton("default-error-handler-tests-1")); config.put("subscriptionName", "default-error-handler-tests-sub-1"); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( - pulsarClient, config); + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, + config); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); PulsarRecordMessageListener messageListener = mock(PulsarRecordMessageListener.class); @@ -112,16 +112,16 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain config.put("topicNames", Collections.singleton("default-error-handler-tests-2")); config.put("subscriptionName", "default-error-handler-tests-sub-2"); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( - pulsarClient, config); + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, + config); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); PulsarRecordMessageListener messageListener = mock(PulsarRecordMessageListener.class); AtomicInteger count = new AtomicInteger(0); doAnswer(invocation -> { - final int currentCount = count.incrementAndGet(); + int currentCount = count.incrementAndGet(); if (currentCount <= 3) { throw new RuntimeException(); } @@ -161,16 +161,16 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain config.put("topicNames", Collections.singleton("default-error-handler-tests-3")); config.put("subscriptionName", "default-error-handler-tests-sub-3"); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( - pulsarClient, config); + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, + config); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); PulsarRecordMessageListener messageListener = mock(PulsarRecordMessageListener.class); doAnswer(invocation -> { - final Message message = invocation.getArgument(1); - final Integer value = message.getValue(); + Message message = invocation.getArgument(1); + Integer value = message.getValue(); if (value % 2 == 0) { throw new RuntimeException(); } @@ -220,26 +220,26 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain config.put("topicNames", Collections.singleton("default-error-handler-tests-4")); config.put("subscriptionName", "default-error-handler-tests-sub-4"); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( - pulsarClient, config); + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, + config); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties.setMaxNumMessages(10); pulsarContainerProperties.setBatchTimeoutMillis(60_000); pulsarContainerProperties.setBatchListener(true); - final PulsarBatchAcknowledgingMessageListener pulsarBatchMessageListener = mock( + PulsarBatchAcknowledgingMessageListener pulsarBatchMessageListener = mock( PulsarBatchAcknowledgingMessageListener.class); doAnswer(invocation -> { - final List> message = invocation.getArgument(1); - final Message integerMessage = message.get(0); - final Integer value = integerMessage.getValue(); + List> message = invocation.getArgument(1); + Message integerMessage = message.get(0); + Integer value = integerMessage.getValue(); if (value == 0) { throw new PulsarBatchListenerFailedException("failed", integerMessage); } - final Acknowledgement acknowledgment = invocation.getArgument(2); + Acknowledgement acknowledgment = invocation.getArgument(2); List messageIds = new ArrayList<>(); for (Message integerMessage1 : message) { messageIds.add(integerMessage1.getMessageId()); @@ -262,9 +262,9 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain Map prodConfig = new HashMap<>(); prodConfig.put("topicName", "default-error-handler-tests-4"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); for (int i = 0; i < 10; i++) { pulsarTemplate.sendAsync(i); } @@ -292,27 +292,27 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain config.put("topicNames", Collections.singleton("default-error-handler-tests-5")); config.put("subscriptionName", "default-error-handler-tests-sub-5"); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( - pulsarClient, config); + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, + config); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties.setMaxNumMessages(10); pulsarContainerProperties.setBatchTimeoutMillis(60_000); pulsarContainerProperties.setBatchListener(true); - final PulsarBatchAcknowledgingMessageListener pulsarBatchMessageListener = mock( + PulsarBatchAcknowledgingMessageListener pulsarBatchMessageListener = mock( PulsarBatchAcknowledgingMessageListener.class); doAnswer(invocation -> { - final List> messages = invocation.getArgument(1); + List> messages = invocation.getArgument(1); for (Message message : messages) { if (message.getValue() == 5) { throw new PulsarBatchListenerFailedException("failed", message); } else { - final Acknowledgement acknowledgment = invocation.getArgument(2); + Acknowledgement acknowledgment = invocation.getArgument(2); acknowledgment.acknowledge(message.getMessageId()); } } @@ -333,9 +333,9 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain Map prodConfig = new HashMap<>(); prodConfig.put("topicName", "default-error-handler-tests-5"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); for (int i = 0; i < 10; i++) { pulsarTemplate.sendAsync(i); } @@ -362,27 +362,27 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain config.put("topicNames", Collections.singleton("default-error-handler-tests-6")); config.put("subscriptionName", "default-error-handler-tests-sub-6"); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( - pulsarClient, config); + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, + config); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties.setMaxNumMessages(10); pulsarContainerProperties.setBatchTimeoutMillis(60_000); pulsarContainerProperties.setBatchListener(true); - final PulsarBatchAcknowledgingMessageListener pulsarBatchMessageListener = mock( + PulsarBatchAcknowledgingMessageListener pulsarBatchMessageListener = mock( PulsarBatchAcknowledgingMessageListener.class); doAnswer(invocation -> { - final List> messages = invocation.getArgument(1); + List> messages = invocation.getArgument(1); for (Message message : messages) { if (message.getValue() == 2 || message.getValue() == 5) { throw new PulsarBatchListenerFailedException("failed", message); } else { - final Acknowledgement acknowledgment = invocation.getArgument(2); + Acknowledgement acknowledgment = invocation.getArgument(2); acknowledgment.acknowledge(message.getMessageId()); } } @@ -403,9 +403,9 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain Map prodConfig = new HashMap<>(); prodConfig.put("topicName", "default-error-handler-tests-6"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); for (int i = 0; i < 10; i++) { pulsarTemplate.sendAsync(i); } @@ -432,25 +432,25 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain config.put("topicNames", Collections.singleton("default-error-handler-tests-7")); config.put("subscriptionName", "default-error-handler-tests-sub-7"); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( - pulsarClient, config); + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, + config); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties.setMaxNumMessages(10); pulsarContainerProperties.setBatchTimeoutMillis(60_000); pulsarContainerProperties.setBatchListener(true); - final PulsarBatchAcknowledgingMessageListener pulsarBatchMessageListener = mock( + PulsarBatchAcknowledgingMessageListener pulsarBatchMessageListener = mock( PulsarBatchAcknowledgingMessageListener.class); AtomicInteger count = new AtomicInteger(0); doAnswer(invocation -> { - final List> messages = invocation.getArgument(1); - final Acknowledgement acknowledgment = invocation.getArgument(2); + List> messages = invocation.getArgument(1); + Acknowledgement acknowledgment = invocation.getArgument(2); for (Message message : messages) { if (message.getValue() == 5) { - final int currentCount = count.getAndIncrement(); + int currentCount = count.getAndIncrement(); if (currentCount < 3) { throw new PulsarBatchListenerFailedException("failed", message); } @@ -479,9 +479,9 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain Map prodConfig = new HashMap<>(); prodConfig.put("topicName", "default-error-handler-tests-7"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); for (int i = 0; i < 10; i++) { pulsarTemplate.sendAsync(i); } @@ -501,25 +501,25 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain config.put("topicNames", Collections.singleton("default-error-handler-tests-8")); config.put("subscriptionName", "default-error-handler-tests-sub-8"); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( - pulsarClient, config); + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, + config); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties.setMaxNumMessages(10); pulsarContainerProperties.setBatchTimeoutMillis(60_000); pulsarContainerProperties.setBatchListener(true); - final PulsarBatchAcknowledgingMessageListener pulsarBatchMessageListener = mock( + PulsarBatchAcknowledgingMessageListener pulsarBatchMessageListener = mock( PulsarBatchAcknowledgingMessageListener.class); AtomicInteger count = new AtomicInteger(0); doAnswer(invocation -> { - final List> messages = invocation.getArgument(1); - final Acknowledgement acknowledgment = invocation.getArgument(2); + List> messages = invocation.getArgument(1); + Acknowledgement acknowledgment = invocation.getArgument(2); for (Message message : messages) { if (message.getValue() == 5) { - final int currentCount = count.getAndIncrement(); + int currentCount = count.getAndIncrement(); if (currentCount < 3) { throw new PulsarBatchListenerFailedException("failed", message); } @@ -551,9 +551,9 @@ public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContain Map prodConfig = new HashMap<>(); prodConfig.put("topicName", "default-error-handler-tests-8"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); for (int i = 0; i < 10; i++) { pulsarTemplate.sendAsync(i); } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java index 0455e2d4..cb8893d7 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java @@ -60,14 +60,14 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS @Test void basicDefaultConsumer() throws Exception { Map config = new HashMap<>(); - final HashSet strings = new HashSet<>(); + HashSet strings = new HashSet<>(); strings.add("dpmlct-012"); config.put("topicNames", strings); config.put("subscriptionName", "dpmlct-sb-012"); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( - pulsarClient, config); + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, + config); CountDownLatch latch = new CountDownLatch(1); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties @@ -78,9 +78,9 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS container.start(); Map prodConfig = new HashMap<>(); prodConfig.put("topicName", "dpmlct-012"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); pulsarTemplate.sendAsync("hello john doe"); assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); container.stop(); @@ -90,15 +90,15 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS @Test void subscriptionInitialPositionEarliest() throws Exception { Map config = new HashMap<>(); - final HashSet strings = new HashSet<>(); + 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(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( - pulsarClient, config); + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, + config); CountDownLatch latch = new CountDownLatch(5); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); pulsarContainerProperties @@ -109,9 +109,9 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS Map prodConfig = new HashMap<>(); prodConfig.put("topicName", "dpmlct-013"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); for (int i = 0; i < 5; i++) { pulsarTemplate.send("hello john doe" + i); } @@ -125,14 +125,14 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS @Test void subscriptionInitialPositionDefaultLatest() throws Exception { Map config = new HashMap<>(); - final HashSet strings = new HashSet<>(); + HashSet strings = new HashSet<>(); strings.add("dpmlct-014"); config.put("topicNames", strings); config.put("subscriptionName", "dpmlct-sb-014"); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( - pulsarClient, config); + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, + config); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); List messages = new ArrayList<>(); pulsarContainerProperties.setMessageListener( @@ -143,9 +143,9 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS Map prodConfig = new HashMap<>(); prodConfig.put("topicName", "dpmlct-014"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); for (int i = 0; i < 5; i++) { pulsarTemplate.send("hello john doe" + i); } @@ -169,9 +169,9 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS .maxDelayMs(5 * 1000).build(); config.put("negativeAckRedeliveryBackoff", redeliveryBackoff); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = spy( + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = spy( new DefaultPulsarConsumerFactory<>(pulsarClient, config)); CountDownLatch latch = new CountDownLatch(10); PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); @@ -186,12 +186,12 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS DefaultPulsarMessageListenerContainer container = new DefaultPulsarMessageListenerContainer<>( pulsarConsumerFactory, pulsarContainerProperties); - final Consumer containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container); + Consumer containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container); Map prodConfig = Collections.singletonMap("topicName", "dpmlct-015"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); for (int i = 0; i < 5; i++) { pulsarTemplate.send("hello john doe" + i); } @@ -219,10 +219,10 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS .deadLetterTopic("dpmlct-016-dlq-topic").build(); config.put("deadLetterPolicy", deadLetterPolicy); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( - pulsarClient, config); + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, + config); CountDownLatch dlqLatch = new CountDownLatch(1); CountDownLatch latch = new CountDownLatch(6); @@ -251,9 +251,9 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS container.start(); Map prodConfig = Collections.singletonMap("topicName", "dpmlct-016"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); for (int i = 1; i < 6; i++) { pulsarTemplate.send(i); } @@ -277,10 +277,10 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS .build(); config.put("deadLetterPolicy", deadLetterPolicy); - final PulsarClient pulsarClient = PulsarClient.builder() - .serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); - final DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>( - pulsarClient, config); + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) + .build(); + DefaultPulsarConsumerFactory pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, + config); CountDownLatch dlqLatch = new CountDownLatch(1); CountDownLatch latch = new CountDownLatch(6); @@ -312,9 +312,9 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS container.start(); Map prodConfig = Collections.singletonMap("topicName", "dpmlct-016"); - final DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>( - pulsarClient, prodConfig); - final PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); + DefaultPulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, + prodConfig); + PulsarTemplate pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory); for (int i = 1; i < 6; i++) { pulsarTemplate.send(i); } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/PulsarListenerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/PulsarListenerTests.java index 9a66794d..16d3068d 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/PulsarListenerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/PulsarListenerTests.java @@ -129,7 +129,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport { @Bean PulsarListenerContainerFactory pulsarListenerContainerFactory( PulsarConsumerFactory pulsarConsumerFactory) { - final ConcurrentPulsarListenerContainerFactory pulsarListenerContainerFactory = new ConcurrentPulsarListenerContainerFactory<>( + ConcurrentPulsarListenerContainerFactory pulsarListenerContainerFactory = new ConcurrentPulsarListenerContainerFactory<>( pulsarConsumerFactory, new PulsarContainerProperties(), null); return pulsarListenerContainerFactory; } @@ -168,9 +168,9 @@ public class PulsarListenerTests implements PulsarTestContainerSupport { void testPulsarListenerProvidedConsumerProperties(@Autowired PulsarListenerEndpointRegistry registry) throws Exception { - final PulsarContainerProperties pulsarContainerProperties = registry.getListenerContainer("foo") + PulsarContainerProperties pulsarContainerProperties = registry.getListenerContainer("foo") .getContainerProperties(); - final Properties pulsarConsumerProperties = pulsarContainerProperties.getPulsarConsumerProperties(); + Properties pulsarConsumerProperties = pulsarContainerProperties.getPulsarConsumerProperties(); assertThat(pulsarConsumerProperties.size()).isEqualTo(2); assertThat(pulsarConsumerProperties.get("topicNames")).isEqualTo("foo-1"); assertThat(pulsarConsumerProperties.get("subscriptionName")).isEqualTo("subscription-1"); @@ -185,7 +185,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport { Map.of("batchingEnabled", false)); PulsarTemplate customTemplate = new PulsarTemplate<>(pulsarProducerFactory); - final ConcurrentPulsarMessageListenerContainer bar = (ConcurrentPulsarMessageListenerContainer) registry + ConcurrentPulsarMessageListenerContainer bar = (ConcurrentPulsarMessageListenerContainer) registry .getListenerContainer("bar"); assertThat(bar.getConcurrency()).isEqualTo(3); @@ -204,7 +204,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport { Map.of("batchingEnabled", false)); PulsarTemplate customTemplate = new PulsarTemplate<>(pulsarProducerFactory); - final ConcurrentPulsarMessageListenerContainer bar = (ConcurrentPulsarMessageListenerContainer) registry + ConcurrentPulsarMessageListenerContainer bar = (ConcurrentPulsarMessageListenerContainer) registry .getListenerContainer("bar"); assertThat(bar.getConcurrency()).isEqualTo(3); @@ -218,7 +218,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport { @Test void ackModeAppliedToContainerFromListener(@Autowired PulsarListenerEndpointRegistry registry) { - final PulsarContainerProperties pulsarContainerProperties = registry.getListenerContainer("ackMode-test-id") + PulsarContainerProperties pulsarContainerProperties = registry.getListenerContainer("ackMode-test-id") .getContainerProperties(); assertThat(pulsarContainerProperties.getAckMode()).isEqualTo(AckMode.RECORD); } @@ -643,7 +643,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport { @Test void simpleListenerWithHeaders() throws Exception { - final MessageId messageId = pulsarTemplate.newMessage("hello-simple-listener") + MessageId messageId = pulsarTemplate.newMessage("hello-simple-listener") .withMessageCustomizer( messageBuilder -> messageBuilder.property("foo", "simpleListenerWithHeaders")) .withTopic("simpleListenerWithHeaders").send(); @@ -657,7 +657,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport { @Test void pulsarMessageListenerWithHeaders() throws Exception { - final MessageId messageId = pulsarTemplate.newMessage("hello-pulsar-message-listener") + MessageId messageId = pulsarTemplate.newMessage("hello-pulsar-message-listener") .withMessageCustomizer( messageBuilder -> messageBuilder.property("foo", "pulsarMessageListenerWithHeaders")) .withTopic("pulsarMessageListenerWithHeaders").send(); @@ -671,7 +671,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport { @Test void springMessagingMessageListenerWithHeaders() throws Exception { - final MessageId messageId = pulsarTemplate.newMessage("hello-spring-messaging-message-listener") + MessageId messageId = pulsarTemplate.newMessage("hello-spring-messaging-message-listener") .withMessageCustomizer(messageBuilder -> messageBuilder.property("foo", "springMessagingMessageListenerWithHeaders")) .withTopic("springMessagingMessageListenerWithHeaders").send(); @@ -685,7 +685,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport { @Test void simpleBatchListenerWithHeaders() throws Exception { - final MessageId messageId = pulsarTemplate.newMessage("hello-simple-batch-listener") + MessageId messageId = pulsarTemplate.newMessage("hello-simple-batch-listener") .withMessageCustomizer( messageBuilder -> messageBuilder.property("foo", "simpleBatchListenerWithHeaders")) .withTopic("simpleBatchListenerWithHeaders").send(); @@ -698,7 +698,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport { @Test void pulsarMessageBatchListenerWithHeaders() throws Exception { - final MessageId messageId = pulsarTemplate.newMessage("hello-pulsar-message-batch-listener") + MessageId messageId = pulsarTemplate.newMessage("hello-pulsar-message-batch-listener") .withMessageCustomizer( messageBuilder -> messageBuilder.property("foo", "pulsarMessageBatchListenerWithHeaders")) .withTopic("pulsarMessageBatchListenerWithHeaders").send(); @@ -712,7 +712,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport { @Test void springMessagingMessageBatchListenerWithHeaders() throws Exception { - final MessageId messageId = pulsarTemplate.newMessage("hello-spring-messaging-message-batch-listener") + MessageId messageId = pulsarTemplate.newMessage("hello-spring-messaging-message-batch-listener") .withMessageCustomizer(messageBuilder -> messageBuilder.property("foo", "springMessagingMessageBatchListenerWithHeaders")) .withTopic("springMessagingMessageBatchListenerWithHeaders").send(); @@ -726,7 +726,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport { @Test void pulsarMessagesBatchListenerWithHeaders() throws Exception { - final MessageId messageId = pulsarTemplate.newMessage("hello-pulsar-messages-batch-listener") + MessageId messageId = pulsarTemplate.newMessage("hello-pulsar-messages-batch-listener") .withMessageCustomizer( messageBuilder -> messageBuilder.property("foo", "pulsarMessagesBatchListenerWithHeaders")) .withTopic("pulsarMessagesBatchListenerWithHeaders").send(); diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/reactive/ReactivePulsarListenerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/reactive/ReactivePulsarListenerTests.java index af036e92..37024110 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/reactive/ReactivePulsarListenerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/reactive/ReactivePulsarListenerTests.java @@ -509,7 +509,7 @@ public class ReactivePulsarListenerTests implements PulsarTestContainerSupport { @Test void simpleListenerWithHeaders() throws Exception { - final MessageId messageId = pulsarTemplate.newMessage("hello-simple-listener") + MessageId messageId = pulsarTemplate.newMessage("hello-simple-listener") .withMessageCustomizer( messageBuilder -> messageBuilder.property("foo", "simpleListenerWithHeaders")) .withTopic("simpleListenerWithHeaders").send(); @@ -523,7 +523,7 @@ public class ReactivePulsarListenerTests implements PulsarTestContainerSupport { @Test void pulsarMessageListenerWithHeaders() throws Exception { - final MessageId messageId = pulsarTemplate.newMessage("hello-pulsar-message-listener") + MessageId messageId = pulsarTemplate.newMessage("hello-pulsar-message-listener") .withMessageCustomizer( messageBuilder -> messageBuilder.property("foo", "pulsarMessageListenerWithHeaders")) .withTopic("pulsarMessageListenerWithHeaders").send(); @@ -537,7 +537,7 @@ public class ReactivePulsarListenerTests implements PulsarTestContainerSupport { @Test void springMessagingMessageListenerWithHeaders() throws Exception { - final MessageId messageId = pulsarTemplate.newMessage("hello-spring-messaging-message-listener") + MessageId messageId = pulsarTemplate.newMessage("hello-spring-messaging-message-listener") .withMessageCustomizer(messageBuilder -> messageBuilder.property("foo", "springMessagingMessageListenerWithHeaders")) .withTopic("springMessagingMessageListenerWithHeaders").send();