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 08be3325..5b8e1bdb 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
@@ -22,6 +22,7 @@ import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;
+import org.apache.pulsar.client.api.MessageId;
import org.apache.pulsar.common.schema.SchemaType;
import org.springframework.messaging.handler.annotation.MessageMapping;
@@ -38,7 +39,7 @@ public @interface PulsarReader {
*
* If none is specified an auto-generated id is used.
*
- * SpEL {@code #{...}} and property place holders {@code ${...}} are supported.
+ * SpEL {@code #{...}} and property placeholders {@code ${...}} are supported.
* @return the {@code id} for the container managing for this endpoint.
* @see PulsarReaderEndpointRegistry#getReaderContainer(String)
*/
@@ -56,6 +57,12 @@ public @interface PulsarReader {
*/
SchemaType schemaType() default SchemaType.NONE;
+ /**
+ * {@link MessageId} for this reader to start from.
+ * @return starting message id - earliest or latest.
+ */
+ String startMessageId() default "";
+
/**
* Topics to listen to.
* @return a comma separated list of topics to listen from.
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 91bbb579..334be910 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
@@ -26,6 +26,8 @@ import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicInteger;
+import org.apache.pulsar.client.api.MessageId;
+
import org.springframework.aop.support.AopUtils;
import org.springframework.beans.BeansException;
import org.springframework.beans.factory.BeanInitializationException;
@@ -210,6 +212,15 @@ public class PulsarReaderAnnotationBeanPostProcessor extends AbstractPulsarAn
endpoint.setId(getEndpointId(pulsarReader));
endpoint.setTopics(topics);
endpoint.setSchemaType(pulsarReader.schemaType());
+ String startMessageIdString = pulsarReader.startMessageId();
+ MessageId startMessageId = null;
+ if (startMessageIdString.equalsIgnoreCase("earliest")) {
+ startMessageId = MessageId.earliest;
+ }
+ else if (startMessageIdString.equalsIgnoreCase("latest")) {
+ startMessageId = MessageId.latest;
+ }
+ endpoint.setStartMessageId(startMessageId);
String autoStartup = pulsarReader.autoStartup();
if (StringUtils.hasText(autoStartup)) {
diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarReaderEndpoint.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarReaderEndpoint.java
index ea6ca65e..4cea0ef9 100644
--- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarReaderEndpoint.java
+++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarReaderEndpoint.java
@@ -21,6 +21,7 @@ import java.util.Arrays;
import java.util.Collections;
import java.util.List;
+import org.apache.pulsar.client.api.MessageId;
import org.apache.pulsar.common.schema.SchemaType;
import org.springframework.beans.BeansException;
@@ -55,6 +56,8 @@ public abstract class AbstractPulsarReaderEndpoint
private final List topics = new ArrayList<>();
+ private MessageId startMessageId;
+
private BeanFactory beanFactory;
private BeanExpressionResolver resolver;
@@ -170,4 +173,12 @@ public abstract class AbstractPulsarReaderEndpoint
this.schemaType = schemaType;
}
+ public MessageId getStartMessageId() {
+ return this.startMessageId;
+ }
+
+ public void setStartMessageId(MessageId startMessageId) {
+ this.startMessageId = startMessageId;
+ }
+
}
diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/DefaultPulsarReaderContainerFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/DefaultPulsarReaderContainerFactory.java
index 51bc3501..4871aebe 100644
--- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/DefaultPulsarReaderContainerFactory.java
+++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/DefaultPulsarReaderContainerFactory.java
@@ -53,6 +53,7 @@ public class DefaultPulsarReaderContainerFactory
}
properties.setSchemaType(endpoint.getSchemaType());
+ properties.setStartMessageId(endpoint.getStartMessageId());
return new DefaultPulsarReaderListenerContainer<>(this.getReaderFactory(), properties);
}
diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarReaderEndpoint.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarReaderEndpoint.java
index f91c02eb..fd6018f3 100644
--- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarReaderEndpoint.java
+++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarReaderEndpoint.java
@@ -18,6 +18,7 @@ package org.springframework.pulsar.config;
import java.util.List;
+import org.apache.pulsar.client.api.MessageId;
import org.apache.pulsar.common.schema.SchemaType;
import org.springframework.lang.Nullable;
@@ -78,4 +79,6 @@ public interface PulsarReaderEndpoint {
@Nullable
Boolean getAutoStartup();
+ MessageId getStartMessageId();
+
}
diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/DefaultPulsarReaderListenerContainer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/DefaultPulsarReaderListenerContainer.java
index 99eb1f74..10084713 100644
--- a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/DefaultPulsarReaderListenerContainer.java
+++ b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/DefaultPulsarReaderListenerContainer.java
@@ -150,7 +150,7 @@ public class DefaultPulsarReaderListenerContainer extends AbstractPulsarReade
readerContainerProperties.getStartMessageId(), (Schema) readerContainerProperties.getSchema());
}
catch (PulsarClientException e) {
- DefaultPulsarReaderListenerContainer.this.logger.error(e, () -> "Pulsar client exceptions.");
+ throw new IllegalStateException("Pulsar client exceptions.", e);
}
}
diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/PulsarReaderContainerProperties.java b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/PulsarReaderContainerProperties.java
index d2b0ab1a..85d70d60 100644
--- a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/PulsarReaderContainerProperties.java
+++ b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/PulsarReaderContainerProperties.java
@@ -25,6 +25,7 @@ import org.apache.pulsar.client.api.Schema;
import org.apache.pulsar.common.schema.SchemaType;
import org.springframework.core.task.AsyncTaskExecutor;
+import org.springframework.pulsar.core.DefaultSchemaResolver;
import org.springframework.pulsar.core.SchemaResolver;
import org.springframework.util.Assert;
@@ -45,7 +46,7 @@ public class PulsarReaderContainerProperties {
private List topics;
- private MessageId startMessageId = MessageId.earliest;
+ private MessageId startMessageId;
private Schema> schema;
@@ -59,6 +60,10 @@ public class PulsarReaderContainerProperties {
return this.readerListener;
}
+ public PulsarReaderContainerProperties() {
+ this.schemaResolver = new DefaultSchemaResolver();
+ }
+
public void setReaderListener(Object readerListener) {
this.readerListener = readerListener;
}
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
new file mode 100644
index 00000000..0f3b7f29
--- /dev/null
+++ b/spring-pulsar/src/test/java/org/springframework/pulsar/reader/PulsarReaderTests.java
@@ -0,0 +1,183 @@
+/*
+ * 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.reader;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+
+import org.apache.pulsar.client.api.PulsarClient;
+import org.junit.jupiter.api.Nested;
+import org.junit.jupiter.api.Test;
+
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.context.annotation.AnnotationConfigApplicationContext;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.pulsar.annotation.EnablePulsar;
+import org.springframework.pulsar.annotation.PulsarReader;
+import org.springframework.pulsar.config.DefaultPulsarReaderContainerFactory;
+import org.springframework.pulsar.config.PulsarClientFactoryBean;
+import org.springframework.pulsar.config.PulsarReaderContainerFactory;
+import org.springframework.pulsar.core.DefaultPulsarProducerFactory;
+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.test.support.PulsarTestContainerSupport;
+import org.springframework.test.annotation.DirtiesContext;
+import org.springframework.test.context.ContextConfiguration;
+import org.springframework.test.context.junit.jupiter.SpringJUnitConfig;
+
+/**
+ * {@link PulsarReader} integration tests.
+ *
+ * @author Soby Chacko
+ */
+@SpringJUnitConfig
+@DirtiesContext
+public class PulsarReaderTests implements PulsarTestContainerSupport {
+
+ @Autowired
+ PulsarTemplate pulsarTemplate;
+
+ @Autowired
+ private PulsarClient pulsarClient;
+
+ @Configuration(proxyBeanMethods = false)
+ @EnablePulsar
+ public static class TopLevelConfig {
+
+ @Bean
+ public PulsarProducerFactory pulsarProducerFactory(PulsarClient pulsarClient) {
+ Map config = Collections.emptyMap();
+ return new DefaultPulsarProducerFactory<>(pulsarClient, config);
+ }
+
+ @Bean
+ public PulsarClientFactoryBean pulsarClientFactoryBean() {
+ return new PulsarClientFactoryBean(Map.of("serviceUrl", PulsarTestContainerSupport.getPulsarBrokerUrl()));
+ }
+
+ @Bean
+ public PulsarTemplate pulsarTemplate(PulsarProducerFactory pulsarProducerFactory) {
+ return new PulsarTemplate<>(pulsarProducerFactory);
+ }
+
+ @Bean
+ public PulsarReaderFactory> pulsarReaderFactory(PulsarClient pulsarClient) {
+ Map config = new HashMap<>();
+ return new DefaultPulsarReaderFactory<>(pulsarClient, config);
+ }
+
+ @Bean
+ PulsarReaderContainerFactory pulsarReaderContainerFactory(PulsarReaderFactory