diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarReader.java b/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarReader.java index 5b8e1bdb..dfff8a65 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarReader.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarReader.java @@ -27,6 +27,7 @@ import org.apache.pulsar.common.schema.SchemaType; import org.springframework.messaging.handler.annotation.MessageMapping; import org.springframework.pulsar.config.PulsarReaderEndpointRegistry; +import org.springframework.pulsar.core.ReaderBuilderCustomizer; @Target({ ElementType.TYPE, ElementType.METHOD, ElementType.ANNOTATION_TYPE }) @Retention(RetentionPolicy.RUNTIME) @@ -75,4 +76,11 @@ public @interface PulsarReader { String autoStartup() default ""; + /** + * The bean name or a 'SpEL' expression that resolves to a + * {@link ReaderBuilderCustomizer} to use to configure the reader. + * @return the bean name or empty string to not configure the reader. + */ + String readerCustomizer() default ""; + } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarReaderAnnotationBeanPostProcessor.java b/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarReaderAnnotationBeanPostProcessor.java index 334be910..b0ffec2a 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarReaderAnnotationBeanPostProcessor.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarReaderAnnotationBeanPostProcessor.java @@ -44,6 +44,8 @@ import org.springframework.pulsar.config.PulsarReaderContainerFactory; import org.springframework.pulsar.config.PulsarReaderEndpoint; import org.springframework.pulsar.config.PulsarReaderEndpointRegistrar; import org.springframework.pulsar.config.PulsarReaderEndpointRegistry; +import org.springframework.pulsar.core.ConsumerBuilderCustomizer; +import org.springframework.pulsar.core.ReaderBuilderCustomizer; import org.springframework.util.Assert; import org.springframework.util.StringUtils; @@ -227,6 +229,23 @@ public class PulsarReaderAnnotationBeanPostProcessor extends AbstractPulsarAn endpoint.setAutoStartup(resolveExpressionAsBoolean(autoStartup, "autoStartup")); } endpoint.setBeanFactory(this.beanFactory); + + resolveReaderCustomizer(endpoint, pulsarReader); + } + + private void resolveReaderCustomizer(MethodPulsarReaderEndpoint endpoint, PulsarReader pulsarReader) { + Object readerCustomizer = resolveExpression(pulsarReader.readerCustomizer()); + if (readerCustomizer instanceof ConsumerBuilderCustomizer) { + endpoint.setReaderBuilderCustomizer((ReaderBuilderCustomizer) readerCustomizer); + } + else { + String readerCustomizerBeanName = resolveExpressionAsString(pulsarReader.readerCustomizer(), + "readerCustomizer"); + if (StringUtils.hasText(readerCustomizerBeanName)) { + endpoint.setReaderBuilderCustomizer( + this.beanFactory.getBean(readerCustomizerBeanName, ReaderBuilderCustomizer.class)); + } + } } private String getEndpointSubscriptionName(PulsarReader pulsarReader) { diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarReaderEndpoint.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarReaderEndpoint.java index a51c802b..bf5b4e00 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarReaderEndpoint.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarReaderEndpoint.java @@ -34,6 +34,7 @@ import org.springframework.messaging.converter.SmartMessageConverter; import org.springframework.messaging.handler.annotation.Header; import org.springframework.messaging.handler.annotation.support.MessageHandlerMethodFactory; import org.springframework.messaging.handler.invocation.InvocableHandlerMethod; +import org.springframework.pulsar.core.ReaderBuilderCustomizer; import org.springframework.pulsar.core.SchemaResolver; import org.springframework.pulsar.listener.Acknowledgement; import org.springframework.pulsar.listener.adapter.AbstractPulsarMessageToSpringMessageAdapter; @@ -63,6 +64,8 @@ public class MethodPulsarReaderEndpoint extends AbstractPulsarReaderEndpoint< private MessageHandlerMethodFactory messageHandlerMethodFactory; + private ReaderBuilderCustomizer readerBuilderCustomizer; + public void setBean(Object bean) { this.bean = bean; } @@ -136,6 +139,8 @@ public class MethodPulsarReaderEndpoint extends AbstractPulsarReaderEndpoint< // })); // } + container.setReaderCustomizer(this.readerBuilderCustomizer); + return readerListener; } @@ -194,4 +199,8 @@ public class MethodPulsarReaderEndpoint extends AbstractPulsarReaderEndpoint< this.messageHandlerMethodFactory = messageHandlerMethodFactory; } + public void setReaderBuilderCustomizer(ReaderBuilderCustomizer readerBuilderCustomizer) { + this.readerBuilderCustomizer = readerBuilderCustomizer; + } + } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultPulsarReaderFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultPulsarReaderFactory.java index f5154be2..fe3f076e 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultPulsarReaderFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/DefaultPulsarReaderFactory.java @@ -53,8 +53,8 @@ public class DefaultPulsarReaderFactory implements PulsarReaderFactory { } @Override - public Reader createReader(@Nullable List topics, @Nullable MessageId messageId, Schema schema) - throws PulsarClientException { + public Reader createReader(@Nullable List topics, @Nullable MessageId messageId, Schema schema, + @Nullable List> customizers) throws PulsarClientException { Objects.requireNonNull(schema, "Schema must be specified"); ReaderBuilder readerBuilder = this.pulsarClient.newReader(schema); if (!CollectionUtils.isEmpty(topics)) { @@ -63,6 +63,11 @@ public class DefaultPulsarReaderFactory implements PulsarReaderFactory { readerBuilder.startMessageId(messageId); readerBuilder.loadConf(this.readerConfig); + + if (!CollectionUtils.isEmpty(customizers)) { + customizers.forEach(customizer -> customizer.customize(readerBuilder)); + } + return readerBuilder.create(); } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarReaderFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarReaderFactory.java index 0df5e4a6..50a01020 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarReaderFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarReaderFactory.java @@ -39,10 +39,13 @@ public interface PulsarReaderFactory { * @param topics set of topics to read from * @param messageId starting message id to read from * @param schema schema of the message to consume + * @param customizers the optional list of customizers to apply to the reader builder. + * Note that the customizers are applied last and have the potential for overriding + * any specified parameters or default properties. * @return Pulsar {@link Reader} * @throws PulsarClientException if there are issues when creating the reader */ - Reader createReader(@Nullable List topics, @Nullable MessageId messageId, Schema schema) - throws PulsarClientException; + Reader createReader(@Nullable List topics, @Nullable MessageId messageId, Schema schema, + @Nullable List> customizers) throws PulsarClientException; } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/ReaderBuilderCustomizer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/ReaderBuilderCustomizer.java new file mode 100644 index 00000000..6d7db0e2 --- /dev/null +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/ReaderBuilderCustomizer.java @@ -0,0 +1,36 @@ +/* + * Copyright 2023 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.pulsar.core; + +import org.apache.pulsar.client.api.ReaderBuilder; + +/** + * The interface to customize a {@link ReaderBuilder}. + * + * @param The message payload type + * @author Soby Chacko + */ +@FunctionalInterface +public interface ReaderBuilderCustomizer { + + /** + * Customizes a {@link ReaderBuilder}. + * @param readerBuilder the builder to customize + */ + void customize(ReaderBuilder readerBuilder); + +} diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/AbstractPulsarMessageReaderContainer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/AbstractPulsarMessageReaderContainer.java index 6b8bb295..78cf8edc 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/AbstractPulsarMessageReaderContainer.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/AbstractPulsarMessageReaderContainer.java @@ -20,6 +20,7 @@ import org.apache.pulsar.client.api.ReaderListener; import org.springframework.pulsar.core.AbstractPulsarMessageContainer; import org.springframework.pulsar.core.PulsarReaderFactory; +import org.springframework.pulsar.core.ReaderBuilderCustomizer; import org.springframework.util.Assert; /** @@ -37,6 +38,8 @@ public non-sealed abstract class AbstractPulsarMessageReaderContainer extends protected final Object lifecycleMonitor = new Object(); + protected ReaderBuilderCustomizer readerBuilderCustomizer; + @SuppressWarnings("unchecked") protected AbstractPulsarMessageReaderContainer(PulsarReaderFactory pulsarReaderFactory, PulsarReaderContainerProperties pulsarReaderContainerProperties) { @@ -96,4 +99,14 @@ public non-sealed abstract class AbstractPulsarMessageReaderContainer extends } } + @SuppressWarnings("unchecked") + @Override + public void setReaderCustomizer(ReaderBuilderCustomizer readerBuilderCustomizer) { + this.readerBuilderCustomizer = (ReaderBuilderCustomizer) readerBuilderCustomizer; + } + + public ReaderBuilderCustomizer getReaderBuilderCustomizer() { + return this.readerBuilderCustomizer; + } + } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/DefaultPulsarMessageReaderContainer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/DefaultPulsarMessageReaderContainer.java index d8872672..4e652915 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/DefaultPulsarMessageReaderContainer.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/DefaultPulsarMessageReaderContainer.java @@ -17,6 +17,8 @@ package org.springframework.pulsar.reader; import java.io.IOException; +import java.util.Collections; +import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; @@ -32,6 +34,7 @@ import org.springframework.context.ApplicationEventPublisher; import org.springframework.core.task.AsyncTaskExecutor; import org.springframework.core.task.SimpleAsyncTaskExecutor; import org.springframework.pulsar.core.PulsarReaderFactory; +import org.springframework.pulsar.core.ReaderBuilderCustomizer; import org.springframework.pulsar.event.ReaderFailedToStartEvent; import org.springframework.pulsar.event.ReaderStartedEvent; import org.springframework.pulsar.event.ReaderStartingEvent; @@ -139,15 +142,21 @@ public class DefaultPulsarMessageReaderContainer extends AbstractPulsarMessag private Reader reader; + private final ReaderBuilderCustomizer readerBuilderCustomizer; + @SuppressWarnings({ "unchecked", "rawtypes" }) InternalAsyncReader(ReaderListener readerListener, PulsarReaderContainerProperties readerContainerProperties) { this.listener = readerListener; this.readerContainerProperties = readerContainerProperties; + this.readerBuilderCustomizer = getReaderBuilderCustomizer(); try { + List> customizers = this.readerBuilderCustomizer != null + ? List.of(this.readerBuilderCustomizer) : Collections.emptyList(); this.reader = getPulsarReaderFactory().createReader(readerContainerProperties.getTopics(), - readerContainerProperties.getStartMessageId(), (Schema) readerContainerProperties.getSchema()); + readerContainerProperties.getStartMessageId(), (Schema) readerContainerProperties.getSchema(), + customizers); } catch (PulsarClientException e) { throw new IllegalStateException("Pulsar client exceptions.", e); diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/PulsarMessageReaderContainer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/PulsarMessageReaderContainer.java index 875c960a..9ef89672 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/PulsarMessageReaderContainer.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/PulsarMessageReaderContainer.java @@ -18,6 +18,7 @@ package org.springframework.pulsar.reader; import org.springframework.beans.factory.DisposableBean; import org.springframework.context.SmartLifecycle; +import org.springframework.pulsar.core.ReaderBuilderCustomizer; /** * Internal abstraction used by the framework representing a message listener container. @@ -43,4 +44,10 @@ public sealed interface PulsarMessageReaderContainer // empty } + /** + * Set a reader customizer on this container. + * @param readerBuilderCustomizer {@link ReaderBuilderCustomizer} + */ + void setReaderCustomizer(ReaderBuilderCustomizer readerBuilderCustomizer); + } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultPulsarReaderFactoryTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultPulsarReaderFactoryTests.java index 78b13ecb..eb59d8da 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultPulsarReaderFactoryTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/DefaultPulsarReaderFactoryTests.java @@ -72,7 +72,7 @@ public class DefaultPulsarReaderFactoryTests implements PulsarTestContainerSuppo void readingFromTheBeginningOfTheTopic() throws Exception { Message message; try (Reader reader = pulsarReaderFactory.createReader(List.of("basic-pulsar-reader-topic"), - MessageId.earliest, Schema.STRING)) { + MessageId.earliest, Schema.STRING, Collections.emptyList())) { Map prodConfig = Map.of("topicName", "basic-pulsar-reader-topic"); PulsarProducerFactory pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, @@ -99,7 +99,7 @@ public class DefaultPulsarReaderFactoryTests implements PulsarTestContainerSuppo Message message; try (Reader reader = pulsarReaderFactory.createReader(List.of("reading-from-the-middle-of-topic"), - messageIds[4], Schema.STRING)) { + messageIds[4], Schema.STRING, Collections.emptyList())) { for (int i = 0; i < 5; i++) { message = reader.readNext(); assertThat(message.getValue()).isEqualTo("hello john doe-" + (i + 5)); @@ -120,7 +120,7 @@ public class DefaultPulsarReaderFactoryTests implements PulsarTestContainerSuppo pulsarTemplate.send("hello john doe"); try (Reader reader = pulsarReaderFactory.createReader(List.of("basic-pulsar-reader-topic"), - MessageId.latest, Schema.STRING)) { + MessageId.latest, Schema.STRING, Collections.emptyList())) { pulsarTemplate.send("hello alice doe"); // It should not read the first message sent (john doe) as latest is the // message id to start. @@ -146,16 +146,15 @@ public class DefaultPulsarReaderFactoryTests implements PulsarTestContainerSuppo @Test void missingTopic() { // topic name is not set in the API call or in the reader config. - assertThatThrownBy( - () -> pulsarReaderFactory.createReader(Collections.emptyList(), MessageId.earliest, Schema.STRING)) - .isInstanceOf(PulsarClientException.class) + assertThatThrownBy(() -> pulsarReaderFactory.createReader(Collections.emptyList(), MessageId.earliest, + Schema.STRING, Collections.emptyList())).isInstanceOf(PulsarClientException.class) .hasMessageContaining("Topic name must be set on the reader builder"); } @Test void missingStartingMessageId() { - assertThatThrownBy(() -> pulsarReaderFactory.createReader(List.of("my-reader-topic"), null, Schema.STRING)) - .isInstanceOf(PulsarClientException.class).hasMessageContaining( + assertThatThrownBy(() -> pulsarReaderFactory.createReader(List.of("my-reader-topic"), null, Schema.STRING, + Collections.emptyList())).isInstanceOf(PulsarClientException.class).hasMessageContaining( "Start message id or start message from roll back must be specified but they cannot be specified at the same time"); } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/reader/PulsarReaderTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/reader/PulsarReaderTests.java index 0f3b7f29..a8475500 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/reader/PulsarReaderTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/reader/PulsarReaderTests.java @@ -25,7 +25,10 @@ import java.util.Map; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; +import org.apache.pulsar.client.api.Message; +import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.PulsarClient; +import org.apache.pulsar.client.api.PulsarClientException; import org.junit.jupiter.api.Nested; import org.junit.jupiter.api.Test; @@ -43,6 +46,7 @@ import org.springframework.pulsar.core.DefaultPulsarReaderFactory; import org.springframework.pulsar.core.PulsarProducerFactory; import org.springframework.pulsar.core.PulsarReaderFactory; import org.springframework.pulsar.core.PulsarTemplate; +import org.springframework.pulsar.core.ReaderBuilderCustomizer; import org.springframework.pulsar.test.support.PulsarTestContainerSupport; import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.context.ContextConfiguration; @@ -180,4 +184,50 @@ public class PulsarReaderTests implements PulsarTestContainerSupport { } + @Nested + @ContextConfiguration(classes = StartMessageIdFromTheMiddleOfTheTopic.WithCustomizerConfig.class) + class StartMessageIdFromTheMiddleOfTheTopic { + + private static final CountDownLatch latch = new CountDownLatch(5); + + @Test + void startMessageIdProvidedThroughReaderCustomizer() throws Exception { + assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); + } + + @EnablePulsar + @Configuration + static class WithCustomizerConfig { + + int currentIndex = 5; + + MessageId[] messageIds = new MessageId[10]; + + @PulsarReader(id = "with-customizer-reader", subscriptionName = "with-customizer-reader-subscription", + topics = "with-customizer-reader-topic", readerCustomizer = "myCustomizer") + void listen(Message message) { + assertThat(message.getMessageId()).isEqualTo(messageIds[currentIndex++]); + latch.countDown(); + } + + @Bean + public ReaderBuilderCustomizer myCustomizer(PulsarTemplate pulsarTemplate) { + return cb -> { + for (int i = 0; i < 10; i++) { + try { + messageIds[i] = pulsarTemplate.send("with-customizer-reader-topic", "hello john doe-"); + } + catch (PulsarClientException e) { + // Ignore + } + } + cb.startMessageId(messageIds[4]); // the first message read is the one + // after this message id. + }; + } + + } + + } + }