From 0048b8aef4050affcf3ed5f966a62dbcde457f10 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Mon, 20 Jul 2020 08:19:04 +0200 Subject: [PATCH] GH-2006 Add test for KafkaNull Resolves #2006 --- spring-cloud-stream/pom.xml | 5 +++ .../SmartMessageMethodArgumentResolver.java | 2 +- .../ImplicitFunctionBindingTests.java | 33 +++++++++++++++++++ 3 files changed, 39 insertions(+), 1 deletion(-) diff --git a/spring-cloud-stream/pom.xml b/spring-cloud-stream/pom.xml index 04ec1410d..c1600bac0 100644 --- a/spring-cloud-stream/pom.xml +++ b/spring-cloud-stream/pom.xml @@ -76,6 +76,11 @@ spring-boot-starter-web test + + org.springframework.kafka + spring-kafka + test + diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/SmartMessageMethodArgumentResolver.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/SmartMessageMethodArgumentResolver.java index 67e2a78b2..187abdd8d 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/SmartMessageMethodArgumentResolver.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/SmartMessageMethodArgumentResolver.java @@ -112,7 +112,7 @@ class SmartMessageMethodArgumentResolver extends MessageMethodArgumentResolver { return !StringUtils.hasText((String) payload); } else { - return false; + return "org.springframework.kafka.support.KafkaNull".equals(payload.getClass().getName()); } } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java index e85189a76..81bb5e0e6 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java @@ -55,6 +55,7 @@ import org.springframework.integration.dsl.IntegrationFlows; import org.springframework.integration.handler.LoggingHandler; import org.springframework.integration.scheduling.PollerMetadata; import org.springframework.integration.support.MessageBuilder; +import org.springframework.kafka.support.KafkaNull; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.support.GenericMessage; @@ -142,6 +143,26 @@ public class ImplicitFunctionBindingTests { } } + @Test + public void testNullMessage() { + + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(NullMessagerConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false", "--spring.cloud.function.definition=func")) { + + InputDestination inputDestination = context.getBean(InputDestination.class); + OutputDestination outputDestination = context.getBean(OutputDestination.class); + + + Message inputMessage = MessageBuilder.withPayload(KafkaNull.INSTANCE).build(); + inputDestination.send(inputMessage); + + Message outputMessage = outputDestination.receive(); + assertThat(outputMessage.getPayload()).isEqualTo("NULL".getBytes()); + } + } + @Test public void testSimpleFunctionWithStreamProperty() { @@ -910,6 +931,18 @@ public class ImplicitFunctionBindingTests { } } + @EnableAutoConfiguration + public static class NullMessagerConfiguration { + + @Bean + public Function func() { + return value -> { + assertThat(value).isNull(); + return "NULL"; + }; + } + } + @EnableAutoConfiguration public static class SingleReactiveConsumerConfiguration {