From 688f05cbc92132fae22873b6d8041205ffff2a1a Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 12 Jul 2019 18:26:31 -0400 Subject: [PATCH] Joining two input KStreams When joining two input KStreams, the binder throws an exceptin. Fixing this issue by passing the wrapped target KStream from the proxy object to the adapted StreamListener method. Resolves #701 --- .../binder/kafka/streams/KStreamBinder.java | 11 ++++- .../streams/KStreamBoundElementFactory.java | 5 +- ...fkaStreamsBindingInformationCatalogue.java | 9 ++++ ...StreamListenerSetupMethodOrchestrator.java | 7 ++- .../StreamToTableJoinIntegrationTests.java | 49 +++++++++++++++++++ 5 files changed, 76 insertions(+), 5 deletions(-) diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java index 4d38188e7..197a6358b 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java @@ -25,6 +25,7 @@ import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.Produced; +import org.springframework.aop.framework.Advised; import org.springframework.cloud.stream.binder.AbstractBinder; import org.springframework.cloud.stream.binder.BinderSpecificPropertiesProvider; import org.springframework.cloud.stream.binder.Binding; @@ -96,8 +97,14 @@ class KStreamBinder extends // @checkstyle:off ExtendedConsumerProperties properties) { // @checkstyle:on - this.kafkaStreamsBindingInformationCatalogue - .registerConsumerProperties(inputTarget, properties.getExtension()); +// this.kafkaStreamsBindingInformationCatalogue +// .registerConsumerProperties(inputTarget, properties.getExtension()); + + KStream delegate = ((KStreamBoundElementFactory.KStreamWrapperHandler) + ((Advised) inputTarget).getAdvisors()[0].getAdvice()).getDelegate(); + + this.kafkaStreamsBindingInformationCatalogue.registerConsumerProperties(delegate, properties.getExtension()); + if (!StringUtils.hasText(group)) { group = this.binderConfigurationProperties.getApplicationId(); } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBoundElementFactory.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBoundElementFactory.java index 9376c53f1..a14436113 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBoundElementFactory.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBoundElementFactory.java @@ -102,7 +102,7 @@ class KStreamBoundElementFactory extends AbstractBindingTargetFactory { } - private static class KStreamWrapperHandler + static class KStreamWrapperHandler implements KStreamWrapper, MethodInterceptor { private KStream delegate; @@ -133,6 +133,9 @@ class KStreamBoundElementFactory extends AbstractBindingTargetFactory { } } + public KStream getDelegate() { + return delegate; + } } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBindingInformationCatalogue.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBindingInformationCatalogue.java index 0115c8454..d055bb0ac 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBindingInformationCatalogue.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBindingInformationCatalogue.java @@ -152,4 +152,13 @@ class KafkaStreamsBindingInformationCatalogue { Serde getKeySerde(KStream kStreamTarget) { return this.keySerdeInfo.get(kStreamTarget); } + + + public Map, BindingProperties> getBindingProperties() { + return bindingProperties; + } + + public Map, KafkaStreamsConsumerProperties> getConsumerProperties() { + return consumerProperties; + } } 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 ea8671fa9..b9dcf78be 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 @@ -273,14 +273,17 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator extends AbstractKafkaStr // wrap the proxy created during the initial target type binding // with real object (KStream) kStreamWrapper.wrap((KStream) stream); - this.kafkaStreamsBindingInformationCatalogue.addKeySerde((KStream) kStreamWrapper, keySerde); + this.kafkaStreamsBindingInformationCatalogue.addKeySerde(stream, keySerde); + BindingProperties bindingProperties1 = this.kafkaStreamsBindingInformationCatalogue.getBindingProperties().get(kStreamWrapper); + this.kafkaStreamsBindingInformationCatalogue.registerBindingProperties(stream, bindingProperties1); + this.kafkaStreamsBindingInformationCatalogue .addStreamBuilderFactory(streamsBuilderFactoryBean); for (StreamListenerParameterAdapter streamListenerParameterAdapter : adapters) { if (streamListenerParameterAdapter.supports(stream.getClass(), methodParameter)) { arguments[parameterIndex] = streamListenerParameterAdapter - .adapt(kStreamWrapper, methodParameter); + .adapt(stream, methodParameter); break; } } 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 7dc0306b1..54157cc0c 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 @@ -20,6 +20,7 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.List; import java.util.Map; +import java.util.concurrent.TimeUnit; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; @@ -32,6 +33,7 @@ import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; import org.apache.kafka.streams.KeyValue; +import org.apache.kafka.streams.kstream.JoinWindows; import org.apache.kafka.streams.kstream.Joined; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.KTable; @@ -350,6 +352,20 @@ public class StreamToTableJoinIntegrationTests { //See this issue: https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/536 } + @Test + public void testTwoKStreamsCanBeJoined() { + SpringApplication app = new SpringApplication( + JoinProcessor.class); + app.setWebApplicationType(WebApplicationType.NONE); + app.run("--server.port=0", + "--spring.cloud.stream.kafka.streams.binder.brokers=" + + embeddedKafka.getBrokersAsString(), + "--spring.application.name=" + + "two-kstream-input-join-integ-test"); + //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/701 + } + @EnableBinding(KafkaStreamsProcessorX.class) @EnableAutoConfiguration @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) @@ -435,4 +451,37 @@ public class StreamToTableJoinIntegrationTests { } + interface BindingsForTwoKStreamJoinTest { + + String INPUT_1 = "input_1"; + String INPUT_2 = "input_2"; + + @Input(INPUT_1) + KStream input_1(); + + @Input(INPUT_2) + KStream input_2(); + } + + @EnableBinding(BindingsForTwoKStreamJoinTest.class) + @EnableAutoConfiguration + public static class JoinProcessor { + + @StreamListener + public void testProcessor( + @Input(BindingsForTwoKStreamJoinTest.INPUT_1) KStream input1Stream, + @Input(BindingsForTwoKStreamJoinTest.INPUT_2) KStream input2Stream) { + input1Stream + .join(input2Stream, + (event1, event2) -> null, + JoinWindows.of(TimeUnit.MINUTES.toMillis(5)), + Joined.with( + Serdes.String(), + Serdes.String(), + Serdes.String() + ) + ); + } + } + }