From 9182adcf560400981f17e7b6bfc01d9d774eee98 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 29 Dec 2017 20:06:00 -0500 Subject: [PATCH] KStream binder outbound keys changes The serializer can fall back to the default common one if there is not a more specific one provided. Resolves #271 --- .../cloud/stream/binder/kstream/KStreamBinder.java | 12 +++++++++++- .../kstream/config/KStreamBinderConfiguration.java | 5 +++-- 2 files changed, 14 insertions(+), 3 deletions(-) diff --git a/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamBinder.java b/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamBinder.java index 41bac7373..7a5253f33 100644 --- a/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamBinder.java +++ b/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamBinder.java @@ -31,6 +31,7 @@ import org.springframework.cloud.stream.binder.DefaultBinding; import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; import org.springframework.cloud.stream.binder.ExtendedProducerProperties; import org.springframework.cloud.stream.binder.ExtendedPropertiesBinder; +import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaProducerProperties; import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner; @@ -54,8 +55,12 @@ public class KStreamBinder extends private final StreamsConfig streamsConfig; - public KStreamBinder(KafkaTopicProvisioner kafkaTopicProvisioner, + private final KafkaBinderConfigurationProperties binderConfigurationProperties; + + public KStreamBinder(KafkaBinderConfigurationProperties binderConfigurationProperties, + KafkaTopicProvisioner kafkaTopicProvisioner, KStreamExtendedBindingProperties kStreamExtendedBindingProperties, StreamsConfig streamsConfig) { + this.binderConfigurationProperties = binderConfigurationProperties; this.kafkaTopicProvisioner = kafkaTopicProvisioner; this.kStreamExtendedBindingProperties = kStreamExtendedBindingProperties; this.streamsConfig = streamsConfig; @@ -94,6 +99,11 @@ public class KStreamBinder extends ((Configurable) keySerde).configure(streamsConfig.originals()); } } + else { + keySerde = this.binderConfigurationProperties.getConfiguration().containsKey("key.serde") ? + Utils.newInstance(this.binderConfigurationProperties.getConfiguration().get("key.serde"), Serde.class) : Serdes.ByteArray(); + } + if (StringUtils.hasText(properties.getExtension().getValueSerde())) { valueSerde = Utils.newInstance(properties.getExtension().getValueSerde(), Serde.class); if (valueSerde instanceof Configurable) { diff --git a/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/config/KStreamBinderConfiguration.java b/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/config/KStreamBinderConfiguration.java index ded4b5169..0db7dd0bb 100644 --- a/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/config/KStreamBinderConfiguration.java +++ b/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/config/KStreamBinderConfiguration.java @@ -49,9 +49,10 @@ public class KStreamBinderConfiguration { } @Bean - public KStreamBinder kStreamBinder(KafkaTopicProvisioner kafkaTopicProvisioner, + public KStreamBinder kStreamBinder(KafkaBinderConfigurationProperties binderConfigurationProperties, + KafkaTopicProvisioner kafkaTopicProvisioner, KStreamExtendedBindingProperties kStreamExtendedBindingProperties, StreamsConfig streamsConfig) { - return new KStreamBinder(kafkaTopicProvisioner, kStreamExtendedBindingProperties, + return new KStreamBinder(binderConfigurationProperties, kafkaTopicProvisioner, kStreamExtendedBindingProperties, streamsConfig); }