diff --git a/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarAnnotationDrivenConfiguration.java b/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarAnnotationDrivenConfiguration.java index c4abe67f..e8a0cac5 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarAnnotationDrivenConfiguration.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarAnnotationDrivenConfiguration.java @@ -61,7 +61,7 @@ public class PulsarAnnotationDrivenConfiguration { map.from(properties::getSchemaType).to(containerProperties::setSchemaType); map.from(properties::getAckMode).to(containerProperties::setAckMode); - map.from(properties::getBatchTimeout).to(containerProperties::setBatchTimeout); + map.from(properties::getBatchTimeoutMillis).to(containerProperties::setBatchTimeout); map.from(properties::getMaxNumBytes).to(containerProperties::setMaxNumBytes); map.from(properties::getMaxNumMessages).to(containerProperties::setMaxNumMessages); diff --git a/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarProperties.java b/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarProperties.java index 3ab6ef5c..a712b937 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarProperties.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarProperties.java @@ -1324,19 +1324,20 @@ public class PulsarProperties { private SchemaType schemaType; /** - * Max number of messages for batch. + * Max number of messages in a single batch request. */ private int maxNumMessages = -1; /** - * Max number of bytes for batch. + * Max number of bytes in a single batch request. */ private int maxNumBytes = 10 * 1024 * 1024; /** - * Timeout for batch. + * Number of milliseconds to wait for enough message to fill a batch request + * before timing out. */ - private int batchTimeout = 100; + private int batchTimeoutMillis = 100; public AckMode getAckMode() { return this.ackMode; @@ -1370,12 +1371,12 @@ public class PulsarProperties { this.maxNumBytes = maxNumBytes; } - public int getBatchTimeout() { - return this.batchTimeout; + public int getBatchTimeoutMillis() { + return this.batchTimeoutMillis; } - public void setBatchTimeout(int batchTimeout) { - this.batchTimeout = batchTimeout; + public void setBatchTimeoutMillis(int batchTimeoutMillis) { + this.batchTimeoutMillis = batchTimeoutMillis; } } diff --git a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarAutoConfigurationTests.java b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarAutoConfigurationTests.java index 43b1758a..8e275d98 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarAutoConfigurationTests.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarAutoConfigurationTests.java @@ -204,8 +204,8 @@ class PulsarAutoConfigurationTests { @Test void consumerBatchPropertiesAreHonored() { contextRunner - .withPropertyValues("spring.pulsar.listener.maxNumMessages=10", - "spring.pulsar.listener.maxNumBytes=101", "spring.pulsar.listener.batchTimeout=50") + .withPropertyValues("spring.pulsar.listener.max-num-messages=10", + "spring.pulsar.listener.max-num-bytes=101", "spring.pulsar.listener.batch-timeout-millis=50") .run((context -> assertThat(context).hasNotFailed() .getBean(ConcurrentPulsarListenerContainerFactory.class).extracting("containerProperties") .hasFieldOrPropertyWithValue("maxNumMessages", 10) 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 0a22a9ea..c2f1a4a9 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 @@ -361,9 +361,9 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess this.nackableMessages.add(message.getMessageId()); } else { - throw new IllegalStateException( - "Exception occurred and not negatively acknowledged or handled properly", - e); + throw new IllegalStateException(String.format( + "Exception occurred and message %s was not auto-nacked; switch to AckMode BATCH or RECORD to enable auto-nacks", + message.getMessageId()), e); } } }