diff --git a/pom.xml b/pom.xml
index 6e293443a..a8d6060fb 100644
--- a/pom.xml
+++ b/pom.xml
@@ -129,6 +129,12 @@
+
+ org.springframework.cloud
+ spring-cloud-stream-schema
+ ${spring-cloud-stream.version}
+ test
+
diff --git a/spring-cloud-stream-binder-kafka-streams/pom.xml b/spring-cloud-stream-binder-kafka-streams/pom.xml
index 4880ddd6f..3f03769ec 100644
--- a/spring-cloud-stream-binder-kafka-streams/pom.xml
+++ b/spring-cloud-stream-binder-kafka-streams/pom.xml
@@ -13,6 +13,10 @@
2.1.0.BUILD-SNAPSHOT
+
+ 1.8.2
+
+
org.springframework.cloud
@@ -76,5 +80,40 @@
testtest
+
+
+ org.springframework.cloud
+ spring-cloud-stream-schema
+ test
+
+
+ org.apache.avro
+ avro
+ ${avro.version}
+ provided
+
+
+
+
+
+ org.apache.avro
+ avro-maven-plugin
+ ${avro.version}
+
+
+ generate-sources
+
+ schema
+ protocol
+ idl-protocol
+
+
+ src/test/resources/avro
+
+
+
+
+
+
\ No newline at end of file
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 ea7d415fa..fb25bb8cc 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
@@ -36,6 +36,7 @@ import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsBinderConfigurationProperties;
import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsExtendedBindingProperties;
+import org.springframework.cloud.stream.binder.kafka.streams.serde.CompositeNonNativeSerde;
import org.springframework.cloud.stream.binding.BindingService;
import org.springframework.cloud.stream.binding.StreamListenerResultAdapter;
import org.springframework.cloud.stream.config.BindingServiceConfiguration;
@@ -161,6 +162,11 @@ public class KafkaStreamsBinderSupportAutoConfiguration {
KafkaStreamsBindingInformationCatalogue, binderConfigurationProperties);
}
+ @Bean
+ public CompositeNonNativeSerde compositeNonNativeSerde(CompositeMessageConverterFactory compositeMessageConverterFactory) {
+ return new CompositeNonNativeSerde(compositeMessageConverterFactory);
+ }
+
@Bean
public KStreamBoundElementFactory kStreamBoundElementFactory(BindingServiceProperties bindingServiceProperties,
KafkaStreamsBindingInformationCatalogue KafkaStreamsBindingInformationCatalogue) {
diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java
index fefef6da4..fddcef7e0 100644
--- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java
+++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java
@@ -21,6 +21,8 @@ import java.util.Map;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
+import org.apache.kafka.common.header.Header;
+import org.apache.kafka.common.header.Headers;
import org.apache.kafka.streams.KeyValue;
import org.apache.kafka.streams.kstream.KStream;
import org.apache.kafka.streams.processor.Processor;
@@ -92,7 +94,7 @@ public class KafkaStreamsMessageConversionDelegate {
}
/**
- * Deserialize incoming {@link KStream} based on contentType.
+ * Deserialize incoming {@link KStream} based on content type.
*
* @param valueClass on KStream value
* @param bindingTarget inbound KStream target
@@ -101,6 +103,9 @@ public class KafkaStreamsMessageConversionDelegate {
@SuppressWarnings("unchecked")
public KStream deserializeOnInbound(Class> valueClass, KStream, ?> bindingTarget) {
MessageConverter messageConverter = compositeMessageConverterFactory.getMessageConverterForAllRegistered();
+ final PerRecordContentTypeHolder perRecordContentTypeHolder = new PerRecordContentTypeHolder();
+
+ resolvePerRecordContentType(bindingTarget, perRecordContentTypeHolder);
//Deserialize using a branching strategy
KStream, ?>[] branch = bindingTarget.branch(
@@ -113,16 +118,26 @@ public class KafkaStreamsMessageConversionDelegate {
if (o2 != null) {
if (valueClass.isAssignableFrom(o2.getClass())) {
keyValueThreadLocal.set(new KeyValue<>(o, o2));
- } else if (o2 instanceof Message) {
- if (valueClass.isAssignableFrom(((Message) o2).getPayload().getClass())) {
- keyValueThreadLocal.set(new KeyValue<>(o, ((Message) o2).getPayload()));
- } else {
- convertAndSetMessage(o, valueClass, messageConverter, (Message) o2);
+ }
+ else if (o2 instanceof Message) {
+ Message> m1 = (Message) o2;
+ if (perRecordContentTypeHolder.contentType != null) {
+ m1 = MessageBuilder.fromMessage(m1).setHeader("contentType", perRecordContentTypeHolder.contentType).build();
}
- } else if (o2 instanceof String || o2 instanceof byte[]) {
- Message> message = MessageBuilder.withPayload(o2).build();
+
+ if (valueClass.isAssignableFrom(m1.getPayload().getClass())) {
+ keyValueThreadLocal.set(new KeyValue<>(o, m1.getPayload()));
+ }
+ else {
+ convertAndSetMessage(o, valueClass, messageConverter, m1);
+ }
+ }
+ else if (o2 instanceof String || o2 instanceof byte[]) {
+ Message> message = perRecordContentTypeHolder.contentType != null ? MessageBuilder.withPayload(o2)
+ .setHeader("contentType", perRecordContentTypeHolder.contentType).build() : MessageBuilder.withPayload(o2).build();
convertAndSetMessage(o, valueClass, messageConverter, message);
- } else {
+ }
+ else {
keyValueThreadLocal.set(new KeyValue<>(o, o2));
}
isValidRecord = true;
@@ -150,6 +165,45 @@ public class KafkaStreamsMessageConversionDelegate {
});
}
+ private static class PerRecordContentTypeHolder {
+
+ String contentType;
+
+ void setContentType(String contentType) {
+ this.contentType = contentType;
+ }
+ }
+
+ @SuppressWarnings("unchecked")
+ private void resolvePerRecordContentType(KStream, ?> outboundBindTarget, PerRecordContentTypeHolder perRecordContentTypeHolder) {
+ outboundBindTarget.process(() -> new Processor() {
+
+ ProcessorContext context;
+
+ @Override
+ public void init(ProcessorContext context) {
+ this.context = context;
+ }
+
+ @Override
+ public void process(Object key, Object value) {
+ final Headers headers = context.headers();
+ final Iterable contentTypes = headers.headers("contentType");
+ if (contentTypes != null && contentTypes.iterator().hasNext()) {
+ final String contentType = new String(contentTypes.iterator().next().value());
+ //remove leading and trailing quotes
+ final String cleanContentType = StringUtils.replace(contentType, "\"", "");
+ perRecordContentTypeHolder.setContentType(cleanContentType);
+ }
+ }
+
+ @Override
+ public void close() {
+
+ }
+ });
+ }
+
private void convertAndSetMessage(Object o, Class> valueClass, MessageConverter messageConverter, Message> msg) {
Object messageConverted = messageConverter.fromMessage(msg, valueClass);
if (messageConverted == null) {
diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CompositeNonNativeSerde.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CompositeNonNativeSerde.java
new file mode 100644
index 000000000..995e97615
--- /dev/null
+++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CompositeNonNativeSerde.java
@@ -0,0 +1,210 @@
+/*
+ * Copyright 2018 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.cloud.stream.binder.kafka.streams.serde;
+
+import java.nio.charset.StandardCharsets;
+import java.util.HashMap;
+import java.util.Map;
+
+import org.apache.kafka.common.serialization.Deserializer;
+import org.apache.kafka.common.serialization.Serde;
+import org.apache.kafka.common.serialization.Serializer;
+
+import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory;
+import org.springframework.messaging.Message;
+import org.springframework.messaging.MessageHeaders;
+import org.springframework.messaging.converter.MessageConverter;
+import org.springframework.messaging.support.MessageBuilder;
+import org.springframework.util.Assert;
+import org.springframework.util.MimeType;
+import org.springframework.util.MimeTypeUtils;
+
+/**
+ * A {@link Serde} implementation that wraps the list of {@link MessageConverter}s
+ * from {@link CompositeMessageConverterFactory}.
+ *
+ * The primary motivation for this class is to provide an avro based {@link Serde} that is
+ * compatible with the schema registry that Spring Cloud Stream provides. When using the
+ * schema registry support from Spring Cloud Stream in a Kafka Streams binder based application,
+ * the applications can deserialize the incoming Kafka Streams records using the built in
+ * Avro {@link MessageConverter}. However, this same message conversion approach will not work
+ * downstream in other operations in the topology for Kafka Streams as some of them needs a
+ * {@link Serde} instance that can talk to the Spring Cloud Stream provided Schema Registry.
+ * This implementation will solve that problem.
+ *
+ * Only Avro and JSON based converters are exposed as binder provided {@link Serde} implementations currently.
+ *
+ * Users of this class must call the {@link CompositeNonNativeSerde#configure(Map, boolean)} method
+ * to configure the {@link Serde} object. At the very least the configuration map must include a key
+ * called "valueClass" to indicate the type of the target object for deserialization. If any other
+ * content type other than JSON is needed (only Avro is available now other than JSON), that needs
+ * to be included in the configuration map with the key "contentType". For example,
+ *
+ *