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() + ) + ); + } + } + }