From 1329655ae6e57ad4a394b924dfdea51bfbe4700f Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Thu, 2 Jul 2020 12:00:22 -0400 Subject: [PATCH] GH-1998: StreamBridge API and conversion logic. When native encoding is enabled, the StreamBridge API still tries to convert using the default converter causing runtime exceptions. Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/1998 --- .../cloud/stream/function/StreamBridge.java | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java index bbd13449c..940de41d5 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java @@ -36,6 +36,7 @@ import org.springframework.cloud.stream.binding.BindingService; import org.springframework.cloud.stream.config.BindingServiceProperties; import org.springframework.cloud.stream.messaging.DirectWithAttributesChannel; import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.integration.support.MessageBuilder; import org.springframework.messaging.Message; import org.springframework.messaging.SubscribableChannel; import org.springframework.util.MimeType; @@ -135,7 +136,11 @@ public final class StreamBridge implements SmartInitializingSingleton { ProducerProperties producerProperties = this.bindingServiceProperties.getProducerProperties(bindingName); SubscribableChannel messageChannel = this.resolveDestination(bindingName, producerProperties); - Function functionToInvoke = this.functionCatalog.lookup(STREAM_BRIDGE_FUNC_NAME, outputContentType.toString()); + boolean skipConversion = producerProperties.isUseNativeEncoding(); + + Function functionToInvoke = skipConversion ? v -> MessageBuilder.withPayload(v).build() : + this.functionCatalog.lookup(STREAM_BRIDGE_FUNC_NAME, outputContentType.toString()); + if (producerProperties != null && producerProperties.isPartitioned()) { functionToInvoke = new PartitionAwareFunctionWrapper((FunctionInvocationWrapper) functionToInvoke, this.applicationContext, producerProperties); }