Kafka streams outbound converter changes

Use getMessageConverterForAllRegistered() from CompositeMessageConverterFactory

Resolves #353
This commit is contained in:
Soby Chacko
2018-04-02 16:41:02 -04:00
committed by Oleg Zhurakousky
parent 84f0fb28ae
commit 2c3787faa1

View File

@@ -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 :