From cf6cea65259b04f959cacf2911ee6987bdcd885d Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 29 Sep 2023 21:24:27 -0400 Subject: [PATCH] GH-2817: Method name clash in Kafka Streams binder - When there are two methods with the same name but with different type erasures, Kafka Streams binder sometimes detects the incorrect method. Fixing this issue by specifically type checking the return type for Kafka Streams types. Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/2817 --- .../kafka/streams/KafkaStreamsBinderUtils.java | 17 ++++++++++++++++- .../MultipleFunctionsInSameAppTests.java | 7 ++++++- 2 files changed, 22 insertions(+), 2 deletions(-) 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));