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 8792a504..c4abe67f 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,6 +61,9 @@ 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::getMaxNumBytes).to(containerProperties::setMaxNumBytes); + map.from(properties::getMaxNumMessages).to(containerProperties::setMaxNumMessages); return factory; } 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 c6cd62fa..3ab6ef5c 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 @@ -1323,6 +1323,21 @@ public class PulsarProperties { */ private SchemaType schemaType; + /** + * Max number of messages for batch. + */ + private int maxNumMessages = -1; + + /** + * Max number of bytes for batch. + */ + private int maxNumBytes = 10 * 1024 * 1024; + + /** + * Timeout for batch. + */ + private int batchTimeout = 100; + public AckMode getAckMode() { return this.ackMode; } @@ -1339,6 +1354,30 @@ public class PulsarProperties { this.schemaType = schemaType; } + public int getMaxNumMessages() { + return this.maxNumMessages; + } + + public void setMaxNumMessages(int maxNumMessages) { + this.maxNumMessages = maxNumMessages; + } + + public int getMaxNumBytes() { + return this.maxNumBytes; + } + + public void setMaxNumBytes(int maxNumBytes) { + this.maxNumBytes = maxNumBytes; + } + + public int getBatchTimeout() { + return this.batchTimeout; + } + + public void setBatchTimeout(int batchTimeout) { + this.batchTimeout = batchTimeout; + } + } public static class Admin { 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 064b7f9f..43b1758a 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 @@ -201,6 +201,18 @@ class PulsarAutoConfigurationTests { InterceptorTestConfiguration.interceptorFoo))); } + @Test + void consumerBatchPropertiesAreHonored() { + contextRunner + .withPropertyValues("spring.pulsar.listener.maxNumMessages=10", + "spring.pulsar.listener.maxNumBytes=101", "spring.pulsar.listener.batchTimeout=50") + .run((context -> assertThat(context).hasNotFailed() + .getBean(ConcurrentPulsarListenerContainerFactory.class).extracting("containerProperties") + .hasFieldOrPropertyWithValue("maxNumMessages", 10) + .hasFieldOrPropertyWithValue("maxNumBytes", 101) + .hasFieldOrPropertyWithValue("batchTimeout", 50))); + } + @Nested class ClientAutoConfigurationTests { 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 46c3faaf..0a22a9ea 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 @@ -360,6 +360,11 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess else if (this.containerProperties.getAckMode() == AckMode.BATCH) { this.nackableMessages.add(message.getMessageId()); } + else { + throw new IllegalStateException( + "Exception occurred and not negatively acknowledged or handled properly", + e); + } } } } @@ -541,12 +546,29 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess @Override public void acknowledge(MessageId messageId) { - throw new UnsupportedOperationException(); + try { + this.consumer.acknowledge(messageId); + } + catch (PulsarClientException e) { + this.consumer.negativeAcknowledge(messageId); + } } @Override public void acknowledge(List messageIds) { - throw new UnsupportedOperationException(); + try { + this.consumer.acknowledge(messageIds); + } + catch (PulsarClientException e) { + for (MessageId messageId : messageIds) { + try { + this.consumer.acknowledge(messageId); + } + catch (PulsarClientException ex) { + this.consumer.negativeAcknowledge(messageId); + } + } + } } @Override @@ -556,7 +578,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess @Override public void nack(MessageId messageId) { - throw new UnsupportedOperationException(); + this.consumer.negativeAcknowledge(messageId); } }