diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java index cde79feaa..9a7f56c11 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java @@ -30,7 +30,9 @@ import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.serialization.ByteArraySerializer; +import org.apache.kafka.streams.kstream.GlobalKTable; import org.apache.kafka.streams.kstream.KStream; +import org.apache.kafka.streams.kstream.KTable; import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.config.BeanDefinition; @@ -49,6 +51,7 @@ import org.springframework.cloud.stream.binder.kafka.utils.DlqDestinationResolve import org.springframework.cloud.stream.binder.kafka.utils.DlqPartitionFunction; import org.springframework.context.ApplicationContext; import org.springframework.core.MethodParameter; +import org.springframework.core.ResolvableType; import org.springframework.kafka.config.StreamsBuilderFactoryBean; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaOperations; @@ -84,7 +87,19 @@ public final class KafkaStreamsBinderUtils { * @return found method as an {@link Optional} */ public static Optional findMethodWithName(String key, Method[] methods) { - return Arrays.stream(methods).filter(m -> m.getName().equals(key)).findFirst(); + return Arrays.stream(methods).filter(m -> m.getName().equals(key) && + returnTypeContainsKafkaStreamsTypes(m)).findFirst(); + } + + private static boolean returnTypeContainsKafkaStreamsTypes(Method method) { + ResolvableType resolvableType = ResolvableType.forMethodReturnType(method); + ResolvableType[] generics = resolvableType.getGenerics(); + if (generics.length > 0) { + Class rawClass = generics[0].getRawClass(); + return rawClass != null && (rawClass.isAssignableFrom(KStream.class) || rawClass.isAssignableFrom(KTable.class) + || rawClass.isAssignableFrom(GlobalKTable.class)); + } + return false; } public static String[] deriveFunctionUnits(String definition) { diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/MultipleFunctionsInSameAppTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/MultipleFunctionsInSameAppTests.java index 22a0efabd..b7d29b245 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/MultipleFunctionsInSameAppTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/MultipleFunctionsInSameAppTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2019-2022 the original author or authors. + * Copyright 2019-2023 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -222,6 +222,11 @@ public class MultipleFunctionsInSameAppTests { (s, p) -> p.equalsIgnoreCase("electronics")); } + // Testing for the scenario under https://github.com/spring-cloud/spring-cloud-stream/issues/2817 + public String processItem(String foo) { + return "testing"; + } + @Bean public Function, KStream> yetAnotherProcess() { return input -> input.map((k, v) -> new KeyValue<>("foo", 1L));