From 024e3cffcc7f2410c52de9b88b2a4b78758155ca Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Tue, 14 Nov 2017 20:29:19 -0500 Subject: [PATCH] Finished consolidating contentType Message conversion - Removed MessageSerializationUtils - Consolidated Message contentType conversion refactoring that was started with d0be34f7cb40be76bdd727b550f9df0292e38040 commit At this point Message conversion is consolidated in either MessageConverters or In/Out channel interceptors configured in MessageConverterConfigurer --- .../cloud/stream/binder/AbstractBinder.java | 7 +- .../binder/AbstractMessageChannelBinder.java | 6 +- .../binder/MessageSerializationUtils.java | 69 ------------------- .../binding/MessageConverterConfigurer.java | 32 +++++++++ 4 files changed, 40 insertions(+), 74 deletions(-) delete mode 100644 spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/MessageSerializationUtils.java diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java index a90e027c6..eb10c3976 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java @@ -146,9 +146,14 @@ public abstract class AbstractBinder message) { - return MessageSerializationUtils.serializePayload(message); + return new MessageValues(message); } @Deprecated diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java index ed8a0fd3c..e26888e7b 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java @@ -16,7 +16,6 @@ package org.springframework.cloud.stream.binder; - import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; @@ -43,7 +42,6 @@ import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.SubscribableChannel; import org.springframework.messaging.support.ChannelInterceptorAdapter; import org.springframework.util.Assert; -import org.springframework.util.MimeType; /** * {@link AbstractBinder} that serves as base class for {@link MessageChannel} binders. @@ -595,11 +593,11 @@ public abstract class AbstractMessageChannelBinder message) { - Object originalPayload = message.getPayload(); - boolean setOriginalContentType = (originalPayload instanceof String); - Assert.isTrue(originalPayload instanceof byte[] || originalPayload instanceof String, - "Failed to convert message's payload. No suitable converter found for provided contentType: " - + message.getHeaders().get(MessageHeaders.CONTENT_TYPE) + " and paylod: " + originalPayload); - Object originalContentType = message.getHeaders().get(MessageHeaders.CONTENT_TYPE); - // Pass content type as String since some transport adapters will exclude - // CONTENT_TYPE Header otherwise - String contentType = null; - if (originalContentType != null) { - contentType = setOriginalContentType ? JavaClassMimeTypeUtils.mimeTypeFromObject(originalPayload, - ObjectUtils.nullSafeToString(originalContentType)).toString() : originalContentType.toString() ; - } - - Object payload = originalPayload instanceof byte[] ? originalPayload : ((String) originalPayload).getBytes(StandardCharsets.UTF_8); - MessageValues messageValues = new MessageValues(message); - messageValues.setPayload(payload); - if (StringUtils.hasText(contentType)) { - messageValues.put(MessageHeaders.CONTENT_TYPE, contentType); - if (originalContentType != null && !originalContentType.toString().equals(contentType.toString())) { - messageValues.put(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE, originalContentType.toString()); - } - } - return messageValues; - } -} diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java index 96f59574e..e3d7601c6 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java @@ -26,6 +26,8 @@ import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; import org.springframework.cloud.stream.binder.BinderException; import org.springframework.cloud.stream.binder.BinderHeaders; import org.springframework.cloud.stream.binder.ConsumerProperties; +import org.springframework.cloud.stream.binder.JavaClassMimeTypeUtils; +import org.springframework.cloud.stream.binder.MessageValues; import org.springframework.cloud.stream.binder.PartitionHandler; import org.springframework.cloud.stream.binder.PartitionKeyExtractorStrategy; import org.springframework.cloud.stream.binder.PartitionSelectorStrategy; @@ -52,6 +54,7 @@ import org.springframework.util.Assert; import org.springframework.util.ClassUtils; import org.springframework.util.MimeType; import org.springframework.util.MimeTypeUtils; +import org.springframework.util.ObjectUtils; import org.springframework.util.StringUtils; /** @@ -354,9 +357,38 @@ public class MessageConverterConfigurer .copyHeaders(headers) .build(); } + postProcessedMessage = this.finishPreSend(postProcessedMessage); } return postProcessedMessage; } + + /** + * This is strictly to support 1.3 semantics where BINDER_ORIGINAL_CONTENT_TYPE header + * needs to be set for certain cases and String payloads needs to be converted to byte[]. + * + * Factored out of what was left of MessageSerializationUtils. + */ + // deprecated at the get go as a reminder to remove at v3.0 + @Deprecated + private Message finishPreSend(Message message) { + String oct = message.getHeaders().containsKey(MessageHeaders.CONTENT_TYPE) ? message.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString() : null; + String ct = oct; + if (message.getPayload() instanceof String) { + ct = JavaClassMimeTypeUtils.mimeTypeFromObject(message.getPayload(), ObjectUtils.nullSafeToString(oct)).toString(); + } + MessageValues messageValues = new MessageValues(message); + Object payload = message.getPayload(); + if (payload instanceof String) { + payload = ((String)payload).getBytes(StandardCharsets.UTF_8); + } + + messageValues.setPayload(payload); + if (ct != null && !ct.equals(oct)) { + messageValues.put(MessageHeaders.CONTENT_TYPE, ct); + messageValues.put(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE, oct); + } + return messageValues.toMessage(); + } }