From 3244fc5274ee4ee5bd6b0fb2bd925b86027fdb41 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Mon, 7 Dec 2020 16:25:00 +0100 Subject: [PATCH] GH-2062 Add logic to log an error for reactive functions Resolves #2062 --- .../stream/function/FunctionConfiguration.java | 3 +++ .../function/ImplicitFunctionBindingTests.java | 17 +++++++++++++++++ 2 files changed, 20 insertions(+) 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 25b66409b..c3e31c5cf 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 @@ -472,6 +472,9 @@ public class FunctionConfiguration { } outputChannel.send((Message) message); } + }).doOnError(e -> { + logger.error("Failure was detected during execution of the reactive function '" + functionDefinition + "'"); + ((Throwable) e).printStackTrace(); }); } if (!function.isConsumer()) { 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 573c1c483..8230732dc 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 @@ -909,6 +909,20 @@ public class ImplicitFunctionBindingTests { } + @Test + public void testGh2062() { + System.clearProperty("spring.cloud.function.definition"); + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(ReactiveFunctionConfiguration.class)) + .web(WebApplicationType.NONE).run("--spring.jmx.enabled=false", + "--spring.cloud.function.definition=echo")) { + InputDestination inputDestination = context.getBean(InputDestination.class); + String jsonPerson = "error"; + + inputDestination.send(MessageBuilder.withPayload(jsonPerson.getBytes()).build()); + // there is really nothing to assert other then check the logs for stack trace and error message + } + } @SuppressWarnings("rawtypes") @Test @@ -1181,6 +1195,9 @@ public class ImplicitFunctionBindingTests { public Function, Flux> echo() { return flux -> flux.map(value -> { System.out.println("echo value reqctive " + value); + if (value.equals("error")) { + throw new RuntimeException("intentional"); + } return value; }); }