From 95a4681d27292ee20985f815cf49468942cd8353 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Mon, 4 Mar 2019 16:20:13 -0500 Subject: [PATCH] KTable input validation KTable as input binding doesn't work if @input is not specified at parameter level. Restructuring the input validation in Kafka Streams binder where it checks for declarative inputs. Fixes #536 --- ...StreamListenerSetupMethodOrchestrator.java | 15 +++++++-- .../StreamToTableJoinIntegrationTests.java | 31 +++++++++++++++++++ 2 files changed, 44 insertions(+), 2 deletions(-) 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. */