diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java index e846781e2..b4926cafd 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java @@ -298,16 +298,20 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro } else if (Function.class.isAssignableFrom(bean.getClass()) || BiFunction.class.isAssignableFrom(bean.getClass())) { Object result; - if (BiFunction.class.isAssignableFrom(bean.getClass())) { - result = ((BiFunction) bean).apply(adaptedInboundArguments[0], adaptedInboundArguments[1]); + + if (composedFunctionNames.length > 0) { + result = handleComposedFunctions(adaptedInboundArguments, null, composedFunctionNames); } else { - result = ((Function) bean).apply(adaptedInboundArguments[0]); - } - result = handleCurriedFunctions(adaptedInboundArguments, result); - if (composedFunctionNames.length > 0) { - result = handleComposedFunctions(adaptedInboundArguments, result, composedFunctionNames); + if (BiFunction.class.isAssignableFrom(bean.getClass())) { + result = ((BiFunction) bean).apply(adaptedInboundArguments[0], adaptedInboundArguments[1]); + } + else { + result = ((Function) bean).apply(adaptedInboundArguments[0]); + } + result = handleCurriedFunctions(adaptedInboundArguments, result); } + if (result != null) { final Set outputs = new TreeSet<>(kafkaStreamsBindableProxyFactory.getOutputs()); final Iterator outboundDefinitionIterator = outputs.iterator();