From 6d4ab483587b1a9ac998830d83ea75be54c1dae2 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Thu, 27 Feb 2020 16:00:36 +0100 Subject: [PATCH] GH-1911 Fix support for Function> Basically the Function> shoudl be treated as Consumer while allowing reference to the Flux so additional operations coudl be applied, but it will no longer result in creation of dummy output destination Resolves #1911 --- .../function/FunctionConfiguration.java | 47 +++++++++++++------ .../ImplicitFunctionBindingTests.java | 28 +++++++++++ 2 files changed, 60 insertions(+), 15 deletions(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java index 323289607..86294b211 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java @@ -18,6 +18,7 @@ package org.springframework.cloud.stream.function; import java.lang.reflect.Field; import java.lang.reflect.Method; +import java.lang.reflect.ParameterizedType; import java.lang.reflect.Type; import java.time.Instant; import java.util.ArrayList; @@ -412,9 +413,12 @@ public class FunctionConfiguration { return MessageChannelReactiveUtils.toPublisher(inputChannel); }).toArray(Publisher[]::new); - BindingProperties bindingProperties = this.serviceProperties.getBindings().get(outputBindingNames.iterator().next()); - ProducerProperties producerProperties = bindingProperties == null ? null : bindingProperties.getProducer(); - PartitionAwareFunction functionToInvoke = new PartitionAwareFunction(function, this.applicationContext, producerProperties); + Function functionToInvoke = function; + if (!CollectionUtils.isEmpty(outputBindingNames)) { + BindingProperties bindingProperties = this.serviceProperties.getBindings().get(outputBindingNames.iterator().next()); + ProducerProperties producerProperties = bindingProperties == null ? null : bindingProperties.getProducer(); + functionToInvoke = new PartitionAwareFunction(function, this.applicationContext, producerProperties); + } Object resultPublishers = functionToInvoke.apply(inputPublishers.length == 1 ? inputPublishers[0] : Tuples.fromArray(inputPublishers)); if (!(resultPublishers instanceof Iterable)) { @@ -422,12 +426,15 @@ public class FunctionConfiguration { } Iterator outputBindingIter = outputBindingNames.iterator(); ((Iterable) resultPublishers).forEach(publisher -> { - MessageChannel outputChannel = this.applicationContext.getBean(outputBindingIter.next(), MessageChannel.class); - Flux.from((Publisher) publisher) - .onErrorContinue((ex, pay) -> { - logger.error("Failed to process the following content which will be dropped: " + pay, (Throwable) ex); - }) - .doOnNext(message -> outputChannel.send((Message) message)).subscribe(); + Flux flux = Flux.from((Publisher) publisher) + .onErrorContinue((ex, pay) -> { + logger.error("Failed to process the following content which will be dropped: " + pay, (Throwable) ex); + }); + if (!CollectionUtils.isEmpty(outputBindingNames)) { + MessageChannel outputChannel = this.applicationContext.getBean(outputBindingIter.next(), MessageChannel.class); + flux = flux.doOnNext(message -> outputChannel.send((Message) message)); + } + flux.subscribe(); }); } else { @@ -437,7 +444,7 @@ public class FunctionConfiguration { Object inputDestination = this.applicationContext.getBean(inputDestinationName); if (inputDestination != null && inputDestination instanceof SubscribableChannel) { ServiceActivatingHandler handler = createFunctionHandler(function, inputDestinationName, outputDestinationName); - if (!FunctionTypeUtils.isConsumer(function.getFunctionType())) { + if (StringUtils.hasText(outputDestinationName)) { // consumer implicit or function<.., mono> handler.setOutputChannelName(outputDestinationName); } ((SubscribableChannel) inputDestination).subscribe(handler); @@ -497,9 +504,7 @@ public class FunctionConfiguration { private boolean isReactiveOrMultipleInputOutput(BindableProxyFactory bindableProxyFactory, Type functionType) { boolean reactiveInputsOutputs = FunctionTypeUtils.isReactive(FunctionTypeUtils.getInputType(functionType, 0)) || FunctionTypeUtils.isReactive(FunctionTypeUtils.getOutputType(functionType, 0)); - return isMultipleInputOutput(bindableProxyFactory) - || (reactiveInputsOutputs - && StringUtils.hasText(this.determineOutputDestinationName(0, bindableProxyFactory, functionType))); + return isMultipleInputOutput(bindableProxyFactory) || reactiveInputsOutputs; } private String determineOutputDestinationName(int index, BindableProxyFactory bindableProxyFactory, Type functionType) { @@ -649,7 +654,7 @@ public class FunctionConfiguration { Type functionType = function.getFunctionType(); if (function.isSupplier()) { this.inputCount = 0; - this.outputCount = FunctionTypeUtils.getOutputCount(functionType); + this.outputCount = this.getOutputCount(functionType, true); } else if (function.isConsumer()) { this.inputCount = FunctionTypeUtils.getInputCount(functionType); @@ -657,7 +662,7 @@ public class FunctionConfiguration { } else { this.inputCount = FunctionTypeUtils.getInputCount(functionType); - this.outputCount = FunctionTypeUtils.getOutputCount(functionType); + this.outputCount = this.getOutputCount(functionType, false); } functionBindableProxyDefinition.getConstructorArgumentValues().addGenericArgumentValue(functionDefinition); @@ -673,6 +678,18 @@ public class FunctionConfiguration { } } + private int getOutputCount(Type functionType, boolean isSupplier) { + int outputCount = FunctionTypeUtils.getOutputCount(functionType); + if (!isSupplier && functionType instanceof ParameterizedType) { + Type outputType = ((ParameterizedType) functionType).getActualTypeArguments()[1]; + if (FunctionTypeUtils.isMono(outputType) && outputType instanceof ParameterizedType + && ((ParameterizedType) outputType).getActualTypeArguments()[0].getTypeName().endsWith("Void")) { + this.outputCount = 0; + } + } + return outputCount; + } + private boolean determineFunctionName(FunctionCatalog catalog, Environment environment) { String definition = streamFunctionProperties.getDefinition(); if (!StringUtils.hasText(definition)) { diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java index 5560c379d..6175a6fd4 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java @@ -30,6 +30,7 @@ import java.util.function.Supplier; import org.junit.After; import org.junit.Test; import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; @@ -677,6 +678,22 @@ public class ImplicitFunctionBindingTests { } } + @Test + public void testReactiveFunctionWithOutputAsMonoVoid() { + System.clearProperty("spring.cloud.function.definition"); + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(FunctionalConsumerConfiguration.class)) + .web(WebApplicationType.NONE).run("--spring.jmx.enabled=false")) { + + InputDestination inputDestination = context.getBean(InputDestination.class); + Message inputMessage = MessageBuilder.withPayload("Hello".getBytes()).build(); + inputDestination.send(inputMessage); + + assertThat(System.getProperty("consumer")).isEqualTo("Hello"); + System.clearProperty("consumer"); + } + } + @EnableAutoConfiguration public static class NoEnableBindingConfiguration { @@ -757,6 +774,17 @@ public class ImplicitFunctionBindingTests { } } + @EnableAutoConfiguration + public static class FunctionalConsumerConfiguration { + @Bean + public Function, Mono> funcConsumer() { + return flux -> flux.doOnNext(value -> { + System.out.println(value); + System.setProperty("consumer", value); + }).then(); + } + } + @EnableAutoConfiguration @EnableBinding(Sink.class) public static class LegacyConfiguration {