From 70cd7dc2f9999dd01ead19a4cc5f5d955187f3ac Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Tue, 1 May 2018 10:03:37 -0400 Subject: [PATCH] Upgrade spring cloud stream version to 2.1.0 sanpshot Always enable multiplex to true in kafka streams binder --- pom.xml | 2 +- .../stream/binder/kafka/streams/KStreamBinder.java | 5 ++++- .../kafka/streams/KStreamBoundElementFactory.java | 4 ++++ .../stream/binder/kafka/streams/KTableBinder.java | 6 +++++- .../kafka/streams/KTableBoundElementFactory.java | 11 ++++++++++- .../KafkaStreamsBinderSupportAutoConfiguration.java | 4 ++-- ...aStreamsStreamListenerSetupMethodOrchestrator.java | 3 ++- 7 files changed, 28 insertions(+), 7 deletions(-) diff --git a/pom.xml b/pom.xml index 7caee65a3..7df879005 100644 --- a/pom.xml +++ b/pom.xml @@ -15,7 +15,7 @@ 2.1.5.RELEASE 3.0.3.RELEASE 1.0.1 - 2.0.0.RELEASE + 2.1.0.BUILD-SNAPSHOT spring-cloud-stream-binder-kafka 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 f4b3250dd..28a26c8f6 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 @@ -92,7 +92,10 @@ class KStreamBinder extends if (!StringUtils.hasText(group)) { group = binderConfigurationProperties.getApplicationId(); } - this.kafkaTopicProvisioner.provisionConsumerDestination(name, group, extendedConsumerProperties); + String[] inputTopics = StringUtils.commaDelimitedListToStringArray(name); + for (String inputTopic : inputTopics) { + this.kafkaTopicProvisioner.provisionConsumerDestination(inputTopic, group, extendedConsumerProperties); + } StreamsConfig streamsConfig = this.KafkaStreamsBindingInformationCatalogue.getStreamsConfig(inputTarget); if (extendedConsumerProperties.getExtension().isEnableDlq()) { String dlqName = StringUtils.isEmpty(extendedConsumerProperties.getExtension().getDlqName()) ? 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 5cb44e718..a45e58a45 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 @@ -21,6 +21,7 @@ import org.aopalliance.intercept.MethodInvocation; import org.apache.kafka.streams.kstream.KStream; import org.springframework.aop.framework.ProxyFactory; +import org.springframework.cloud.stream.binder.ConsumerProperties; import org.springframework.cloud.stream.binding.AbstractBindingTargetFactory; import org.springframework.cloud.stream.config.BindingProperties; import org.springframework.cloud.stream.config.BindingServiceProperties; @@ -50,6 +51,9 @@ class KStreamBoundElementFactory extends AbstractBindingTargetFactory { @Override public KStream createInput(String name) { + ConsumerProperties consumerProperties = this.bindingServiceProperties.getConsumerProperties(name); + //Always set multiplex to true in the kafka streams binder + consumerProperties.setMultiplex(true); return createProxyForKStream(name); } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java index fd9399dc3..a2624bb6c 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java @@ -73,7 +73,11 @@ class KTableBinder extends if (!StringUtils.hasText(group)) { group = binderConfigurationProperties.getApplicationId(); } - this.kafkaTopicProvisioner.provisionConsumerDestination(name, group, extendedConsumerProperties); + + String[] inputTopics = StringUtils.commaDelimitedListToStringArray(name); + for (String inputTopic : inputTopics) { + this.kafkaTopicProvisioner.provisionConsumerDestination(inputTopic, group, extendedConsumerProperties); + } if (extendedConsumerProperties.getExtension().isEnableDlq()) { String dlqName = StringUtils.isEmpty(extendedConsumerProperties.getExtension().getDlqName()) ? diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBoundElementFactory.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBoundElementFactory.java index 915235b71..e1d64fe75 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBoundElementFactory.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBoundElementFactory.java @@ -21,7 +21,9 @@ import org.aopalliance.intercept.MethodInvocation; import org.apache.kafka.streams.kstream.KTable; import org.springframework.aop.framework.ProxyFactory; +import org.springframework.cloud.stream.binder.ConsumerProperties; import org.springframework.cloud.stream.binding.AbstractBindingTargetFactory; +import org.springframework.cloud.stream.config.BindingServiceProperties; import org.springframework.util.Assert; /** @@ -33,12 +35,19 @@ import org.springframework.util.Assert; */ class KTableBoundElementFactory extends AbstractBindingTargetFactory { - KTableBoundElementFactory() { + private final BindingServiceProperties bindingServiceProperties; + + KTableBoundElementFactory(BindingServiceProperties bindingServiceProperties) { super(KTable.class); + this.bindingServiceProperties = bindingServiceProperties; } @Override public KTable createInput(String name) { + ConsumerProperties consumerProperties = this.bindingServiceProperties.getConsumerProperties(name); + //Always set multiplex to true in the kafka streams binder + consumerProperties.setMultiplex(true); + KTableBoundElementFactory.KTableWrapperHandler wrapper= new KTableBoundElementFactory.KTableWrapperHandler(); ProxyFactory proxyFactory = new ProxyFactory(KTableBoundElementFactory.KTableWrapper.class, KTable.class); proxyFactory.addAdvice(wrapper); diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java index 95d0ee290..4ec89519c 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java @@ -122,8 +122,8 @@ public class KafkaStreamsBinderSupportAutoConfiguration { } @Bean - public KTableBoundElementFactory kTableBoundElementFactory() { - return new KTableBoundElementFactory(); + public KTableBoundElementFactory kTableBoundElementFactory(BindingServiceProperties bindingServiceProperties) { + return new KTableBoundElementFactory(bindingServiceProperties); } @Bean 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 9cdd0458c..0a5a4ed9d 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 @@ -333,7 +333,8 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListene } - private KStream getkStream(String inboundName, KafkaStreamsStateStoreProperties storeSpec, BindingProperties bindingProperties, StreamsBuilder streamsBuilder, + private KStream getkStream(String inboundName, KafkaStreamsStateStoreProperties storeSpec, + BindingProperties bindingProperties, StreamsBuilder streamsBuilder, Serde keySerde, Serde valueSerde) { if (storeSpec != null) { StoreBuilder storeBuilder = buildStateStore(storeSpec);