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 1299b8b45..29ef31c98 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 @@ -1,5 +1,5 @@ /* - * Copyright 2018-2021 the original author or authors. + * Copyright 2018-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. @@ -16,9 +16,12 @@ package org.springframework.cloud.stream.binder.kafka.streams; +import java.lang.reflect.Method; +import java.util.Arrays; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.function.BiFunction; import org.apache.commons.logging.Log; @@ -64,7 +67,7 @@ import org.springframework.util.StringUtils; * @author Soby Chacko * @author Gary Russell */ -final class KafkaStreamsBinderUtils { +public final class KafkaStreamsBinderUtils { private static final Log LOGGER = LogFactory.getLog(KafkaStreamsBinderUtils.class); @@ -72,6 +75,17 @@ final class KafkaStreamsBinderUtils { } + /** + * Utility method to find the method targeted by the key. + * + * @param key name of the method + * @param methods collection of methods to search from + * @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(); + } + static void prepareConsumerBinding(String name, String group, ApplicationContext context, KafkaTopicProvisioner kafkaTopicProvisioner, KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/FunctionDetectorCondition.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/FunctionDetectorCondition.java index eac3bd7ac..fc10220a3 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/FunctionDetectorCondition.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/FunctionDetectorCondition.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. @@ -37,6 +37,7 @@ import org.springframework.beans.factory.annotation.AnnotatedBeanDefinition; import org.springframework.beans.factory.config.BeanDefinition; import org.springframework.boot.autoconfigure.condition.ConditionOutcome; import org.springframework.boot.autoconfigure.condition.SpringBootCondition; +import org.springframework.cloud.stream.binder.kafka.streams.KafkaStreamsBinderUtils; import org.springframework.context.annotation.ConditionContext; import org.springframework.core.ResolvableType; import org.springframework.core.type.AnnotatedTypeMetadata; @@ -94,12 +95,17 @@ public class FunctionDetectorCondition extends SpringBootCondition { try { Method[] methods = classObj.getDeclaredMethods(); - Optional kafkaStreamMethod = Arrays.stream(methods).filter(m -> m.getName().equals(key)).findFirst(); - // check if the bean name is overridden. + Optional kafkaStreamMethod = KafkaStreamsBinderUtils.findMethodWithName(key, methods); if (kafkaStreamMethod.isEmpty()) { + // check any inherited methods + methods = classObj.getMethods(); + kafkaStreamMethod = KafkaStreamsBinderUtils.findMethodWithName(key, methods); + } + if (kafkaStreamMethod.isEmpty()) { + // check if the bean name is overridden. final BeanDefinition beanDefinition = context.getBeanFactory().getBeanDefinition(key); final String factoryMethodName = beanDefinition.getFactoryMethodName(); - kafkaStreamMethod = Arrays.stream(methods).filter(m -> m.getName().equals(factoryMethodName)).findFirst(); + kafkaStreamMethod = KafkaStreamsBinderUtils.findMethodWithName(factoryMethodName, methods); } if (kafkaStreamMethod.isPresent()) { diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionBeanPostProcessor.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionBeanPostProcessor.java index 368e4a6c4..655b06168 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionBeanPostProcessor.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionBeanPostProcessor.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. @@ -48,6 +48,7 @@ import org.springframework.beans.factory.config.BeanDefinition; import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; import org.springframework.beans.factory.support.BeanDefinitionRegistry; import org.springframework.beans.factory.support.RootBeanDefinition; +import org.springframework.cloud.stream.binder.kafka.streams.KafkaStreamsBinderUtils; import org.springframework.cloud.stream.function.StreamFunctionProperties; import org.springframework.core.ResolvableType; import org.springframework.util.ClassUtils; @@ -176,11 +177,15 @@ public class KafkaStreamsFunctionBeanPostProcessor implements InitializingBean, ClassUtils.getDefaultClassLoader()); try { Method[] methods = classObj.getDeclaredMethods(); - Optional functionalBeanMethods = Arrays.stream(methods).filter(m -> m.getName().equals(key)).findFirst(); + Optional functionalBeanMethods = KafkaStreamsBinderUtils.findMethodWithName(key, methods); + if (functionalBeanMethods.isEmpty()) { + methods = classObj.getMethods(); // check the inherited methods + functionalBeanMethods = KafkaStreamsBinderUtils.findMethodWithName(key, methods); + } if (functionalBeanMethods.isEmpty()) { final BeanDefinition beanDefinition = this.beanFactory.getBeanDefinition(key); final String factoryMethodName = beanDefinition.getFactoryMethodName(); - functionalBeanMethods = Arrays.stream(methods).filter(m -> m.getName().equals(factoryMethodName)).findFirst(); + functionalBeanMethods = KafkaStreamsBinderUtils.findMethodWithName(factoryMethodName, methods); } if (functionalBeanMethods.isPresent()) { diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/bootstrap/KafkaStreamsBinderBootstrapTest.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/bootstrap/KafkaStreamsBinderBootstrapTest.java index 6ed8b0aac..4746c7c28 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/bootstrap/KafkaStreamsBinderBootstrapTest.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/bootstrap/KafkaStreamsBinderBootstrapTest.java @@ -1,5 +1,5 @@ /* - * Copyright 2018-2022 the original author or authors. + * Copyright 2018-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. @@ -176,7 +176,7 @@ class KafkaStreamsBinderBootstrapTest { } @SpringBootApplication - static class SimpleKafkaStreamsApplication { + static class SimpleKafkaStreamsApplication extends BaseConfig { @Bean public Consumer> input1() { @@ -192,6 +192,11 @@ class KafkaStreamsBinderBootstrapTest { }; } + } + + // Testing the scenario reported by https://github.com/spring-cloud/spring-cloud-stream/issues/2737 + static class BaseConfig { + @Bean public Consumer> input3() { return s -> { @@ -199,4 +204,6 @@ class KafkaStreamsBinderBootstrapTest { }; } } + + }