From 0e2954c01edf7b9995d56b173ff3149369d2ba10 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 29 May 2024 11:54:05 -0400 Subject: [PATCH] GH-2941: Kafka Streams binder composition issues * When composed functions are used on component functions in the Kafka Streams binder, there is an issue in which the first function in the composition is invoked twice. Fixing this issue by ensuring that the function execution path is only invoked once. Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/2941 --- .../streams/KafkaStreamsFunctionProcessor.java | 18 +++++++++++------- 1 file changed, 11 insertions(+), 7 deletions(-) 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();