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 272d139f8..9d4823123 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 @@ -62,8 +62,11 @@ public class KStreamBinder extends private final StreamsConfig streamsConfig; + private final KafkaBinderConfigurationProperties binderConfigurationProperties; + public KStreamBinder(KafkaBinderConfigurationProperties binderConfigurationProperties, KafkaTopicProvisioner kafkaTopicProvisioner, - KStreamExtendedBindingProperties kStreamExtendedBindingProperties, StreamsConfig streamsConfig) { + KStreamExtendedBindingProperties kStreamExtendedBindingProperties, StreamsConfig streamsConfig) { + this.binderConfigurationProperties = binderConfigurationProperties; this.headers = EmbeddedHeaderUtils.headersToEmbed(binderConfigurationProperties.getHeaders()); this.kafkaTopicProvisioner = kafkaTopicProvisioner; this.kStreamExtendedBindingProperties = kStreamExtendedBindingProperties; @@ -72,7 +75,7 @@ public class KStreamBinder extends @Override protected Binding> doBindConsumer(String name, String group, - KStream inputTarget, ExtendedConsumerProperties properties) { + KStream inputTarget, ExtendedConsumerProperties properties) { ExtendedConsumerProperties extendedConsumerProperties = new ExtendedConsumerProperties( new KafkaConsumerProperties()); @@ -83,17 +86,17 @@ public class KStreamBinder extends @Override @SuppressWarnings("unchecked") protected Binding> doBindProducer(String name, KStream outboundBindTarget, - ExtendedProducerProperties properties) { + ExtendedProducerProperties properties) { ExtendedProducerProperties extendedProducerProperties = new ExtendedProducerProperties( new KafkaProducerProperties()); - this.kafkaTopicProvisioner.provisionProducerDestination(name , extendedProducerProperties); + this.kafkaTopicProvisioner.provisionProducerDestination(name, extendedProducerProperties); if (HeaderMode.embeddedHeaders.equals(properties.getHeaderMode())) { outboundBindTarget = outboundBindTarget.map(new KeyValueMapper>() { @Override public KeyValue apply(Object k, Object v) { if (v instanceof Message) { try { - return new KeyValue<>(k, (Object)KStreamBinder.this.serializeAndEmbedHeadersIfApplicable((Message) v)); + return new KeyValue<>(k, (Object) KStreamBinder.this.serializeAndEmbedHeadersIfApplicable((Message) v)); } catch (Exception e) { throw new IllegalArgumentException(e); @@ -111,7 +114,7 @@ public class KStreamBinder extends .map(new KeyValueMapper>() { @Override public KeyValue apply(Object k, Object v) { - return KeyValue.pair(k, (Object)KStreamBinder.this.serializePayloadIfNecessary((Message) v)); + return KeyValue.pair(k, (Object) KStreamBinder.this.serializePayloadIfNecessary((Message) v)); } }); } @@ -126,31 +129,37 @@ public class KStreamBinder extends } } if (!properties.isUseNativeEncoding() || StringUtils.hasText(properties.getExtension().getKeySerde()) || StringUtils.hasText(properties.getExtension().getValueSerde())) { - Serde keySerde = Serdes.ByteArray(); - Serde valueSerde = Serdes.ByteArray(); try { + Serde keySerde; + Serde valueSerde; + if (StringUtils.hasText(properties.getExtension().getKeySerde())) { keySerde = Utils.newInstance(properties.getExtension().getKeySerde(), Serde.class); if (keySerde instanceof Configurable) { ((Configurable) keySerde).configure(streamsConfig.originals()); } } - } - catch (ClassNotFoundException e) { - throw new IllegalStateException("Serde class not found: ", e); - } - try { + 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) { ((Configurable) valueSerde).configure(streamsConfig.originals()); } } + else { + valueSerde = this.binderConfigurationProperties.getConfiguration().containsKey("value.serde") ? + Utils.newInstance(this.binderConfigurationProperties.getConfiguration().get("value.serde"), Serde.class) : Serdes.ByteArray(); + } + outboundBindTarget.to((Serde) keySerde, (Serde) valueSerde, name); } catch (ClassNotFoundException e) { throw new IllegalStateException("Serde class not found: ", e); } - outboundBindTarget.to((Serde) keySerde, (Serde) valueSerde, name); + } else { outboundBindTarget.to(name); diff --git a/spring-cloud-stream-binder-kstream/src/test/java/org/springframework/cloud/stream/binder/kstream/KstreamBinderPojoInputStringOutputIntegrationTests.java b/spring-cloud-stream-binder-kstream/src/test/java/org/springframework/cloud/stream/binder/kstream/KstreamBinderPojoInputStringOutputIntegrationTests.java index 95da3b434..b1838ed4c 100644 --- a/spring-cloud-stream-binder-kstream/src/test/java/org/springframework/cloud/stream/binder/kstream/KstreamBinderPojoInputStringOutputIntegrationTests.java +++ b/spring-cloud-stream-binder-kstream/src/test/java/org/springframework/cloud/stream/binder/kstream/KstreamBinderPojoInputStringOutputIntegrationTests.java @@ -85,6 +85,7 @@ public class KstreamBinderPojoInputStringOutputIntegrationTests { "--spring.cloud.stream.kstream.binder.configuration.value.serde=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.bindings.output.producer.headerMode=raw", "--spring.cloud.stream.bindings.output.producer.useNativeEncoding=true", + "--spring.cloud.stream.kstream.bindings.output.producer.keySerde=org.apache.kafka.common.serialization.Serdes$IntegerSerde", "--spring.cloud.stream.bindings.input.consumer.headerMode=raw", "--spring.cloud.stream.kstream.binder.brokers=" + embeddedKafka.getBrokersAsString(), "--spring.cloud.stream.kstream.binder.zkNodes=" + embeddedKafka.getZookeeperConnectionString()); @@ -108,7 +109,7 @@ public class KstreamBinderPojoInputStringOutputIntegrationTests { @StreamListener("input") @SendTo("output") - public KStream process(KStream input) { + public KStream process(KStream input) { return input .filter(new Predicate() { @@ -128,11 +129,11 @@ public class KstreamBinderPojoInputStringOutputIntegrationTests { .groupByKey(new JsonSerde<>(Product.class), new JsonSerde<>(Product.class)) .count(TimeWindows.of(5000), "id-count-store") .toStream() - .map(new KeyValueMapper, Long, KeyValue>() { + .map(new KeyValueMapper, Long, KeyValue>() { @Override - public KeyValue apply(Windowed key, Long value) { - return new KeyValue<>(null, "Count for product with ID 123: " + value); + public KeyValue apply(Windowed key, Long value) { + return new KeyValue<>(key.key().id, "Count for product with ID 123: " + value); } }); }