diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java index ddda09e96..59b402fb8 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java @@ -629,8 +629,19 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator if (!methodParameter.getParameterType().isAssignableFrom(Object.class) && this.applicationContext.containsBean(targetBeanName)) { Class targetBeanClass = this.applicationContext.getType(targetBeanName); - return this.streamListenerParameterAdapter.supports(targetBeanClass, - methodParameter); + if (targetBeanClass != null) { + boolean supports = KStream.class.isAssignableFrom(targetBeanClass) + && KStream.class.isAssignableFrom(methodParameter.getParameterType()); + if (!supports) { + supports = KTable.class.isAssignableFrom(targetBeanClass) + && KTable.class.isAssignableFrom(methodParameter.getParameterType()); + if (!supports) { + supports = GlobalKTable.class.isAssignableFrom(targetBeanClass) + && GlobalKTable.class.isAssignableFrom(methodParameter.getParameterType()); + } + } + return supports; + } } return false; } diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToTableJoinIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToTableJoinIntegrationTests.java index 05b44bee9..715b632ce 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToTableJoinIntegrationTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToTableJoinIntegrationTests.java @@ -343,6 +343,20 @@ public class StreamToTableJoinIntegrationTests { } } + @Test + public void testTrivialSingleKTableInputAsNonDeclarative() { + SpringApplication app = new SpringApplication( + TrivialKTableApp.class); + app.setWebApplicationType(WebApplicationType.NONE); + app.run("--server.port=0", + "--spring.cloud.stream.kafka.streams.binder.brokers=" + + embeddedKafka.getBrokersAsString(), + "--spring.cloud.stream.kafka.streams.bindings.input-y.consumer.application-id=" + + "testTrivialSingleKTableInputAsNonDeclarative"); + //All we are verifying is that this application didn't throw any errors. + //See this issue: https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/536 + } + @EnableBinding(KafkaStreamsProcessorX.class) @EnableAutoConfiguration @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) @@ -368,6 +382,16 @@ public class StreamToTableJoinIntegrationTests { } + @EnableBinding(KafkaStreamsProcessorY.class) + @EnableAutoConfiguration + public static class TrivialKTableApp { + + @StreamListener("input-y") + public void process(KTable inputTable) { + inputTable.toStream().foreach((key, value) -> System.out.println("key : value " + key + " : " + value)); + } + } + interface KafkaStreamsProcessorX extends KafkaStreamsProcessor { @Input("input-x") @@ -375,6 +399,13 @@ public class StreamToTableJoinIntegrationTests { } + interface KafkaStreamsProcessorY { + + @Input("input-y") + KTable inputY(); + + } + /** * Tuple for a region and its associated number of clicks. */