diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java index 5c767227..b507d742 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java @@ -394,17 +394,14 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess */ private List> invokeBatchListenerErrorHandler(AtomicBoolean inRetryMode, AtomicBoolean messagesPendingInBatch, List> messageList, Throwable exception) { - try { - Assert.isInstanceOf(PulsarBatchListenerFailedException.class, exception, - "Batch listener should throw PulsarBatchListenerFailedException on errors."); - } - catch (Exception e) { - // try in the cause if something downstream wrapped the original - // exception. + + // Make sure either the exception or the exception cause is batch exception + if (!(exception instanceof PulsarBatchListenerFailedException)) { exception = exception.getCause(); Assert.isInstanceOf(PulsarBatchListenerFailedException.class, exception, "Batch listener should throw PulsarBatchListenerFailedException on errors."); } + PulsarBatchListenerFailedException pulsarBatchListenerFailedException = (PulsarBatchListenerFailedException) exception; Message pulsarMessage = getPulsarMessageCausedTheException(pulsarBatchListenerFailedException); final Message theCurrentPulsarMessageTracked = this.pulsarConsumerErrorHandler.currentMessage(); diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/PulsarDeadLetterPublishingRecoverer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/PulsarDeadLetterPublishingRecoverer.java index e202de54..c3ebd11d 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/PulsarDeadLetterPublishingRecoverer.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/PulsarDeadLetterPublishingRecoverer.java @@ -68,8 +68,8 @@ public class PulsarDeadLetterPublishingRecoverer implements PulsarMessageReco this.pulsarTemplate.newMessage(message.getValue()) .withTopic(this.destinationResolver.apply(consumer, message)) .withMessageCustomizer(messageBuilder -> messageBuilder.property(EXCEPTION_THROWN_CAUSE, - exception.getCause() == null ? exception.getMessage() - : exception.getCause().getMessage())) + exception.getCause() != null ? exception.getCause().getMessage() + : exception.getMessage())) .sendAsync(); } catch (PulsarClientException e) {