diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index c06e4a958..f4fe0f877 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -1228,7 +1228,7 @@ public class KafkaMessageChannelBinder extends if (message.getPayload() instanceof Throwable throwablePayload) { throwable = throwablePayload; - + String exceptionMessage = buildMessage(throwable, throwable.getCause()); HeaderMode headerMode = properties.getHeaderMode(); if (headerMode == null || HeaderMode.headers.equals(headerMode)) { @@ -1248,7 +1248,6 @@ public class KafkaMessageChannelBinder extends .getBytes(StandardCharsets.UTF_8))); kafkaHeaders.add(new RecordHeader(X_EXCEPTION_FQCN, throwable .getClass().getName().getBytes(StandardCharsets.UTF_8))); - String exceptionMessage = throwable.getMessage(); if (exceptionMessage != null) { kafkaHeaders.add(new RecordHeader(X_EXCEPTION_MESSAGE, exceptionMessage.getBytes(StandardCharsets.UTF_8))); @@ -1271,8 +1270,7 @@ public class KafkaMessageChannelBinder extends record.timestampType().toString()); messageValues.put(X_EXCEPTION_FQCN, throwable.getClass().getName()); - messageValues.put(X_EXCEPTION_MESSAGE, - throwable.getMessage()); + messageValues.put(X_EXCEPTION_MESSAGE, exceptionMessage); messageValues.put(X_EXCEPTION_STACKTRACE, getStackTraceAsString(throwable)); @@ -1323,6 +1321,26 @@ public class KafkaMessageChannelBinder extends } } + @Nullable + private String buildMessage(Throwable exception, Throwable cause) { + String message = exception.getMessage(); + if (!exception.equals(cause)) { + if (message != null) { + message = message + "; "; + } + String causeMsg = cause.getMessage(); + if (causeMsg != null) { + if (message != null) { + message = message + causeMsg; + } + else { + message = causeMsg; + } + } + } + return message; + } + @SuppressWarnings("unchecked") @Nullable private KafkaAwareTransactionManager transactionManager(@Nullable String beanName) { diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index bb44daa9f..701af2102 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -1300,7 +1300,7 @@ public class KafkaBinderTests extends else { assertThat(new String((byte[]) receivedMessage.getHeaders() .get(KafkaMessageChannelBinder.X_EXCEPTION_MESSAGE))).startsWith( - "Dispatcher failed to deliver Message"); + "Dispatcher failed to deliver Message; fail"); } assertThat(receivedMessage.getHeaders()