From 685472b9966918ef707ac5a7f29285ecfe95ffb9 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 16 Sep 2022 20:45:44 -0400 Subject: [PATCH] Fix batch payload type as List> * List> was not casting properly for batch listeners. Addressing this issue. --- .../config/MethodPulsarListenerEndpoint.java | 3 +++ .../DefaultPulsarMessageListenerContainer.java | 15 ++++++++++++--- .../PulsarDeadLetterPublishingRecoverer.java | 3 ++- ...ulsarBatchMessagingMessageListenerAdapter.java | 2 +- .../PulsarMessagingMessageListenerAdapter.java | 3 ++- 5 files changed, 20 insertions(+), 6 deletions(-) diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarListenerEndpoint.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarListenerEndpoint.java index 161f4a2c..00ea537b 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarListenerEndpoint.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarListenerEndpoint.java @@ -222,6 +222,9 @@ public class MethodPulsarListenerEndpoint extends AbstractPulsarListenerEndpo if (rawClass != null && isContainerType(rawClass)) { resolvableType = resolvableType.getGeneric(0); } + if (Message.class.isAssignableFrom(resolvableType.getRawClass())) { + resolvableType = resolvableType.getGeneric(0); + } return resolvableType; } 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 f2d2259c..5c767227 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 @@ -393,9 +393,18 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess * @return a list of messages to be processed next. */ private List> invokeBatchListenerErrorHandler(AtomicBoolean inRetryMode, - AtomicBoolean messagesPendingInBatch, List> messageList, Exception exception) { - Assert.isInstanceOf(PulsarBatchListenerFailedException.class, exception, - "Batch listener should throw PulsarBatchListenerFailedException on errors."); + 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. + 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 628a79a1..e202de54 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,7 +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().getMessage())) + exception.getCause() == null ? exception.getMessage() + : exception.getCause().getMessage())) .sendAsync(); } catch (PulsarClientException e) { diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/adapter/PulsarBatchMessagingMessageListenerAdapter.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/adapter/PulsarBatchMessagingMessageListenerAdapter.java index 82f2a278..d7322197 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/adapter/PulsarBatchMessagingMessageListenerAdapter.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/adapter/PulsarBatchMessagingMessageListenerAdapter.java @@ -82,7 +82,7 @@ public class PulsarBatchMessagingMessageListenerAdapter extends PulsarMessagi } } else { - message = null; // optimization since we won't need any conversion to invoke + message = MessageBuilder.withPayload(msg).build(); } logger.debug(() -> "Processing [" + message + "]"); diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/adapter/PulsarMessagingMessageListenerAdapter.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/adapter/PulsarMessagingMessageListenerAdapter.java index b020e187..70bc241a 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/adapter/PulsarMessagingMessageListenerAdapter.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/adapter/PulsarMessagingMessageListenerAdapter.java @@ -227,7 +227,8 @@ public abstract class PulsarMessagingMessageListenerAdapter { && parameterizedType.getActualTypeArguments().length == 1) { Type paramType = parameterizedType.getActualTypeArguments()[0]; - this.isConsumerRecordList = paramType.equals(Messages.class); + this.isConsumerRecordList = paramType instanceof ParameterizedType + && ((ParameterizedType) paramType).getRawType().equals(Message.class); boolean messageHasGeneric = paramType instanceof ParameterizedType && ((ParameterizedType) paramType) .getRawType().equals(org.springframework.messaging.Message.class); this.isMessageList = paramType.equals(org.springframework.messaging.Message.class) || messageHasGeneric;