From 3e9fc5152e446438bd52c7c9d028bccbacb42c8d Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Wed, 16 Mar 2022 14:38:33 +0100 Subject: [PATCH] GH-2280 Add initial support for per-binding poller configuration --- .../stream/binder/ProducerProperties.java | 54 +++++++++++++++++++ .../function/FunctionConfiguration.java | 42 +++++++++++++-- .../stream/binding/ExplicitBindingTests.java | 14 ++++- 3 files changed, 106 insertions(+), 4 deletions(-) diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/ProducerProperties.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/ProducerProperties.java index a8a726353..996a61ba7 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/ProducerProperties.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/ProducerProperties.java @@ -17,6 +17,7 @@ package org.springframework.cloud.stream.binder; import java.io.IOException; +import java.time.Duration; import com.fasterxml.jackson.annotation.JsonInclude; import com.fasterxml.jackson.annotation.JsonInclude.Include; @@ -74,6 +75,8 @@ public class ProducerProperties { private boolean errorChannelEnabled = false; + private PollerProperties poller; + public Expression getPartitionKeyExpression() { return this.partitionKeyExpression; } @@ -168,6 +171,14 @@ public class ProducerProperties { this.autoStartup = autoStartup; } + public PollerProperties getPoller() { + return poller; + } + + public void setPoller(PollerProperties poller) { + this.poller = poller; + } + static class ExpressionSerializer extends JsonSerializer { @Override @@ -180,4 +191,47 @@ public class ProducerProperties { } + public static class PollerProperties { + + private Duration fixedDelay = Duration.ofMillis(1000); + + private long maxMessagesPerPoll = 1L; + + private String cron; + + private Duration initialDelay = Duration.ofMillis(0); + + public long getMaxMessagesPerPoll() { + return maxMessagesPerPoll; + } + + public void setMaxMessagesPerPoll(long maxMessagesPerPoll) { + this.maxMessagesPerPoll = maxMessagesPerPoll; + } + + public String getCron() { + return cron; + } + + public void setCron(String cron) { + this.cron = cron; + } + + public Duration getInitialDelay() { + return initialDelay; + } + + public void setInitialDelay(Duration initialDelay) { + this.initialDelay = initialDelay; + } + + public Duration getFixedDelay() { + return fixedDelay; + } + + public void setFixedDelay(Duration fixedDelay) { + this.fixedDelay = fixedDelay; + } + } + } diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java index ffdca9a3a..6ba050b38 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java @@ -21,6 +21,7 @@ import java.lang.reflect.GenericArrayType; import java.lang.reflect.Method; import java.lang.reflect.ParameterizedType; import java.lang.reflect.Type; +import java.time.Duration; import java.time.Instant; import java.util.ArrayList; import java.util.Arrays; @@ -52,6 +53,7 @@ import org.springframework.boot.autoconfigure.AutoConfigureAfter; import org.springframework.boot.autoconfigure.AutoConfigureBefore; import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.boot.context.properties.PropertyMapper; import org.springframework.cloud.function.cloudevent.CloudEventMessageUtils; import org.springframework.cloud.function.context.FunctionCatalog; import org.springframework.cloud.function.context.FunctionProperties; @@ -67,6 +69,7 @@ import org.springframework.cloud.stream.binder.BinderFactory; import org.springframework.cloud.stream.binder.BindingCreatedEvent; import org.springframework.cloud.stream.binder.ConsumerProperties; import org.springframework.cloud.stream.binder.ProducerProperties; +import org.springframework.cloud.stream.binder.ProducerProperties.PollerProperties; import org.springframework.cloud.stream.binding.BindableProxyFactory; import org.springframework.cloud.stream.binding.NewDestinationBindingCallback; import org.springframework.cloud.stream.binding.SupportedBindableFeatures; @@ -95,6 +98,7 @@ import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.IntegrationFlowBuilder; import org.springframework.integration.dsl.IntegrationFlows; import org.springframework.integration.handler.AbstractMessageHandler; +import org.springframework.integration.scheduling.PollerMetadata; import org.springframework.integration.support.MessageBuilder; import org.springframework.integration.util.IntegrationReactiveUtils; import org.springframework.lang.Nullable; @@ -104,6 +108,9 @@ import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.MessagingException; import org.springframework.messaging.SubscribableChannel; import org.springframework.scheduling.TaskScheduler; +import org.springframework.scheduling.Trigger; +import org.springframework.scheduling.support.CronTrigger; +import org.springframework.scheduling.support.PeriodicTrigger; import org.springframework.util.Assert; import org.springframework.util.ClassUtils; import org.springframework.util.CollectionUtils; @@ -313,11 +320,25 @@ public class FunctionConfiguration { taskScheduler.schedule(() -> { }, Instant.now()); // will keep AC alive } else { // implies pollable + AtomicReference pollerMetadata = new AtomicReference<>(); + if (producerProperties != null && producerProperties.getPoller() != null) { + PollerProperties poller = producerProperties.getPoller(); + + PollerMetadata pm = new PollerMetadata(); + PropertyMapper map = PropertyMapper.get().alwaysApplyingWhenNonNull(); + map.from(poller::getMaxMessagesPerPoll).to(pm::setMaxMessagesPerPoll); + map.from(poller).as(this::asTrigger).to(pm::setTrigger); + pollerMetadata.set(pm); + } boolean autoStartup = producerProperties != null ? producerProperties.isAutoStartup() : true; - integrationFlowBuilder = IntegrationFlows - .fromSupplier(supplier, spca -> spca.id(bindingName + "_spca").autoStartup(autoStartup)); - //only apply the PollableBean attributes if this is a reactive function. + integrationFlowBuilder = pollerMetadata == null + ? IntegrationFlows.fromSupplier(supplier, + spca -> spca.id(bindingName + "_spca").autoStartup(autoStartup)) + : IntegrationFlows.fromSupplier(supplier, spca -> spca.id(bindingName + "_spca") + .poller(pollerMetadata.get()).autoStartup(autoStartup)); + + // only apply the PollableBean attributes if this is a reactive function. if (splittable && reactive) { integrationFlowBuilder = integrationFlowBuilder.split(); } @@ -326,6 +347,21 @@ public class FunctionConfiguration { return integrationFlowBuilder; } + private Trigger asTrigger(PollerProperties poller) { + if (StringUtils.hasText(poller.getCron())) { + return new CronTrigger(poller.getCron()); + } + return createPeriodicTrigger(poller.getFixedDelay(), poller.getInitialDelay()); + } + + private Trigger createPeriodicTrigger(Duration period, Duration initialDelay) { + PeriodicTrigger trigger = new PeriodicTrigger(period.toMillis()); + if (initialDelay != null) { + trigger.setInitialDelay(initialDelay.toMillis()); + } + return trigger; + } + private PollableBean extractPollableAnnotation(StreamFunctionProperties functionProperties, GenericApplicationContext context, BindableFunctionProxyFactory proxyFactory) { // here we need to ensure that for cases where composition is defined we only look for supplier method to find Pollable annotation. diff --git a/core/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/ExplicitBindingTests.java b/core/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/ExplicitBindingTests.java index 009c65628..61518d949 100644 --- a/core/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/ExplicitBindingTests.java +++ b/core/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/ExplicitBindingTests.java @@ -24,6 +24,7 @@ import org.junit.jupiter.api.Test; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.builder.SpringApplicationBuilder; +import org.springframework.cloud.stream.binder.test.OutputDestination; import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; @@ -73,10 +74,21 @@ public class ExplicitBindingTests { TestChannelBinderConfiguration.getCompleteConfiguration(SupplierConfiguration.class)) .web(WebApplicationType.NONE) .run("--spring.jmx.enabled=false", - "--spring.cloud.stream.input-bindings=supply")) { + "--spring.cloud.stream.input-bindings=supply", + "--spring.cloud.stream.bindings.supply-out-0.producer.poller.fixed-delay=3000")) { + assertThat(context.getBean("supply-in-0", MessageChannel.class)).isNotNull(); assertThat(context.getBean("supply-out-0", MessageChannel.class)).isNotNull(); + + OutputDestination output = context.getBean(OutputDestination.class); + assertThat(output.receive()).isNotNull(); + assertThat(output.receive(500)).isNull(); + assertThat(output.receive(500)).isNull(); + assertThat(output.receive(500)).isNull(); + assertThat(output.receive(500)).isNull(); + assertThat(output.receive(500)).isNull(); + assertThat(output.receive(500)).isNotNull(); } }