diff --git a/dependencies.gradle b/dependencies.gradle index 622b318a..0ced6561 100644 --- a/dependencies.gradle +++ b/dependencies.gradle @@ -1,5 +1,5 @@ ext { - springBootVersion = '3.2.6' + springBootVersion = '3.2.7' springCloudVersion = '2023.0.2' springCloudAwsVersion = '3.0.4' diff --git a/function/spring-aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/AggregatorFunctionConfiguration.java b/function/spring-aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/AggregatorFunctionConfiguration.java index 032b4596..381e2717 100644 --- a/function/spring-aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/AggregatorFunctionConfiguration.java +++ b/function/spring-aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/AggregatorFunctionConfiguration.java @@ -19,6 +19,7 @@ package org.springframework.cloud.fn.aggregator; import java.util.function.Function; import reactor.core.publisher.Flux; +import reactor.core.scheduler.Schedulers; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.BeanFactoryAware; @@ -67,8 +68,9 @@ public class AggregatorFunctionConfiguration { @Bean public Function>, Flux>> aggregatorFunction(FluxMessageChannel aggregatorInputChannel) { return (input) -> Flux.from(this.outputChannel) - .doOnRequest((request) -> aggregatorInputChannel.subscribeTo(input.map(( - inputMessage) -> MessageBuilder.fromMessage(inputMessage).removeHeader("kafka_consumer").build()))); + .doOnRequest((request) -> aggregatorInputChannel.subscribeTo(input + .map((inputMessage) -> MessageBuilder.fromMessage(inputMessage).removeHeader("kafka_consumer").build()) + .publishOn(Schedulers.boundedElastic()))); } @Bean