diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java index d71d953be..7fc0420ef 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java @@ -557,6 +557,10 @@ public class FunctionConfiguration { @Override public void handleMessageInternal(Message message) throws MessagingException { Object result = functionInvocationWrapper.apply((Message) message); + if (result == null) { + logger.debug("Function execution resulted in null. No message will be sent"); + return; + } if (result instanceof Iterable) { for (Object resultElement : (Iterable) result) { this.doSendMessage(resultElement, message); diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ScenarioTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ScenarioTests.java index 5eeb17b3b..65428a80e 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ScenarioTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ScenarioTests.java @@ -24,6 +24,7 @@ import org.junit.jupiter.api.Test; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.builder.SpringApplicationBuilder; +import org.springframework.cloud.stream.binder.test.InputDestination; import org.springframework.cloud.stream.binder.test.OutputDestination; import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration; import org.springframework.context.ConfigurableApplicationContext; @@ -56,6 +57,23 @@ public class ScenarioTests { } } + @Test + public void test2107() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(SupplierReturningNullConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false", + "--spring.cloud.function.definition=uppercase")) { + + InputDestination input = context.getBean(InputDestination.class); + input.send(new GenericMessage("a".getBytes()), "uppercase-in-0"); + OutputDestination output = context.getBean(OutputDestination.class); + assertThat(new String(output.receive(2000, "uppercase-out-0").getPayload())).isEqualTo("a"); + input.send(new GenericMessage("b".getBytes()), "uppercase-in-0"); + assertThat(output.receive(2000, "uppercase-out-0")).isNull(); + } + } + @EnableAutoConfiguration @Configuration public static class TestConfiguration { @@ -68,4 +86,20 @@ public class ScenarioTests { return message -> message; } } + + @EnableAutoConfiguration + @Configuration + public static class SupplierReturningNullConfiguration { + @Bean + public Function uppercase() { + return v -> { + if ("a".equals(v)) { + return v; + } + else { + return null; + } + }; + } + } }