From 2c3787faa1449dc5411956df46f323dd7a6cbab2 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Mon, 2 Apr 2018 16:41:02 -0400 Subject: [PATCH] Kafka streams outbound converter changes Use getMessageConverterForAllRegistered() from CompositeMessageConverterFactory Resolves #353 --- .../kafka/streams/KafkaStreamsMessageConversionDelegate.java | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) 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 87d98607f..b507fc2cf 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 @@ -30,7 +30,6 @@ 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.MimeType; import org.springframework.util.StringUtils; /** @@ -72,9 +71,7 @@ class KafkaStreamsMessageConversionDelegate { */ public KStream serializeOnOutbound(KStream outboundBindTarget) { String contentType = this.kstreamBindingInformationCatalogue.getContentType(outboundBindTarget); - MessageConverter messageConverter = StringUtils.hasText(contentType) ? compositeMessageConverterFactory - .getMessageConverterForType(MimeType.valueOf(contentType)) - : null; + MessageConverter messageConverter = compositeMessageConverterFactory.getMessageConverterForAllRegistered(); return outboundBindTarget.map((k, v) -> { Message message = v instanceof Message ? (Message) v :