diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/AbstractMethodAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/AbstractMethodAnnotationPostProcessor.java index 31ceb36771..da4d33ebdf 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/AbstractMethodAnnotationPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/AbstractMethodAnnotationPostProcessor.java @@ -688,7 +688,14 @@ public abstract class AbstractMethodAnnotationPostProcessor "No poller has been defined for channel-adapter '" + this.beanName + "', and no default poller is available within the context."); } - if (this.pollerMetadata.getMaxMessagesPerPoll() == Integer.MIN_VALUE) { + long maxMessagesPerPoll = this.pollerMetadata.getMaxMessagesPerPoll(); + if (maxMessagesPerPoll == PollerMetadata.MAX_MESSAGES_UNBOUNDED) { // the default is 1 since a source might return // a non-null and non-interruptible value every time it is invoked - this.pollerMetadata.setMaxMessagesPerPoll(1); + maxMessagesPerPoll = 1; } - spca.setMaxMessagesPerPoll(this.pollerMetadata.getMaxMessagesPerPoll()); + spca.setMaxMessagesPerPoll(maxMessagesPerPoll); if (this.sendTimeout != null) { spca.setSendTimeout(this.sendTimeout); } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java b/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java index 3990919ef4..6505577dea 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java @@ -104,6 +104,7 @@ import org.springframework.integration.endpoint.EventDrivenConsumer; import org.springframework.integration.endpoint.MethodInvokingMessageSource; import org.springframework.integration.endpoint.PollingConsumer; import org.springframework.integration.endpoint.ReactiveStreamsConsumer; +import org.springframework.integration.endpoint.SourcePollingChannelAdapter; import org.springframework.integration.expression.SpelPropertyAccessorRegistrar; import org.springframework.integration.gateway.GatewayProxyFactoryBean; import org.springframework.integration.handler.ServiceActivatingHandler; @@ -416,12 +417,15 @@ public class EnableIntegrationTests { assertThat(this.counterChannel.receive(10)).isNull(); - SmartLifecycle countSA = this.context.getBean("annotationTestService.count.inboundChannelAdapter", - SmartLifecycle.class); + SourcePollingChannelAdapter countSA = + this.context.getBean("annotationTestService.count.inboundChannelAdapter", + SourcePollingChannelAdapter.class); assertThat(countSA.isAutoStartup()).isFalse(); assertThat(countSA.getPhase()).isEqualTo(23); countSA.start(); + assertThat(countSA.getMaxMessagesPerPoll()).isEqualTo(1); + for (int i = 0; i < 10; i++) { Message message = this.counterChannel.receive(10_000); assertThat(message).isNotNull();