From b85ee57d80835956541c182521d0dffd0ef2ce4a Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Thu, 29 Sep 2022 17:50:31 -0400 Subject: [PATCH] Kafka Streams Binder Default Package Beans Currently, in Kafka Streams binder-based apps, processor beans need to be declared public. This is unnecessary and caused by some restrictions in the binder. This PR fixes this restriction. Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/2516 --- .../kafka/streams/function/FunctionDetectorCondition.java | 2 +- .../function/KafkaStreamsFunctionBeanPostProcessor.java | 2 +- .../KafkaStreamsBinderWordCountFunctionTests.java | 8 ++++---- 3 files changed, 6 insertions(+), 6 deletions(-) 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 7c8f1bf59..eac3bd7ac 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 @@ -93,7 +93,7 @@ public class FunctionDetectorCondition extends SpringBootCondition { ClassUtils.getDefaultClassLoader()); try { - Method[] methods = classObj.getMethods(); + Method[] methods = classObj.getDeclaredMethods(); Optional kafkaStreamMethod = Arrays.stream(methods).filter(m -> m.getName().equals(key)).findFirst(); // check if the bean name is overridden. if (kafkaStreamMethod.isEmpty()) { 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 829f8a655..368e4a6c4 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 @@ -175,7 +175,7 @@ public class KafkaStreamsFunctionBeanPostProcessor implements InitializingBean, .getMetadata().getClassName(), ClassUtils.getDefaultClassLoader()); try { - Method[] methods = classObj.getMethods(); + Method[] methods = classObj.getDeclaredMethods(); Optional functionalBeanMethods = Arrays.stream(methods).filter(m -> m.getName().equals(key)).findFirst(); if (functionalBeanMethods.isEmpty()) { final BeanDefinition beanDefinition = this.beanFactory.getBeanDefinition(key); diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java index 1a61ddf31..53f12d9bb 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java @@ -391,7 +391,7 @@ public class KafkaStreamsBinderWordCountFunctionTests { InteractiveQueryService interactiveQueryService; @Bean - public Function, KStream> process() { + Function, KStream> process() { return input -> input .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+"))) @@ -405,7 +405,7 @@ public class KafkaStreamsBinderWordCountFunctionTests { } @Bean - public StreamsBuilderFactoryBeanConfigurer customizer() { + StreamsBuilderFactoryBeanConfigurer customizer() { return fb -> { try { fb.setStateListener((newState, oldState) -> { @@ -422,7 +422,7 @@ public class KafkaStreamsBinderWordCountFunctionTests { } @Bean - public StreamPartitioner streamPartitioner() { + StreamPartitioner streamPartitioner() { return (t, k, v, n) -> k.equals("foo") ? 0 : 1; } } @@ -431,7 +431,7 @@ public class KafkaStreamsBinderWordCountFunctionTests { static class OutboundNullApplication { @Bean - public Function, KStream> process() { + Function, KStream> process() { return input -> input .flatMapValues( value -> Arrays.asList(value.toLowerCase().split("\\W+")))