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 b507d742..da2de243 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 @@ -170,6 +170,12 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess private final PulsarConsumerErrorHandler pulsarConsumerErrorHandler; + private final boolean isBatchListener = this.containerProperties.isBatchListener(); + + private final AckMode ackMode = this.containerProperties.getAckMode(); + + private final SubscriptionType subscriptionType = this.containerProperties.getSubscriptionType(); + @SuppressWarnings({ "unchecked", "rawtypes" }) Listener(MessageListener messageListener) { if (messageListener instanceof PulsarBatchMessageListener) { @@ -282,7 +288,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess DefaultPulsarMessageListenerContainer.this.logger.error(e, () -> "Error receiving messages."); } Assert.isTrue(messages != null, "Messages cannot be null."); - if (this.containerProperties.isBatchListener()) { + if (this.isBatchListener) { if (!inRetryMode.get() && !messagesPendingInBatch.get()) { messageList = new ArrayList<>(); messages.forEach(messageList::add); @@ -291,13 +297,13 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess if (messageList != null && messageList.size() > 0) { if (this.batchMessageListener instanceof PulsarBatchAcknowledgingMessageListener) { this.batchMessageListener.received(this.consumer, messageList, - this.containerProperties.getAckMode() == AckMode.MANUAL + this.ackMode.equals(AckMode.MANUAL) ? new ConsumerBatchAcknowledgment(this.consumer) : null); } else { this.batchMessageListener.received(this.consumer, messageList); } - if (this.containerProperties.getAckMode() == AckMode.BATCH) { + if (this.ackMode.equals(AckMode.BATCH)) { try { if (isSharedSubscriptionType()) { this.consumer.acknowledge(messages); @@ -335,29 +341,26 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess do { try { if (this.listener instanceof PulsarAcknowledgingMessageListener) { - this.listener.received(this.consumer, message, - this.containerProperties.getAckMode() == AckMode.MANUAL - ? new ConsumerAcknowledgment(this.consumer, message) : null); + this.listener.received(this.consumer, message, this.ackMode.equals(AckMode.MANUAL) + ? new ConsumerAcknowledgment(this.consumer, message) : null); } else if (this.listener != null) { this.listener.received(this.consumer, message); } - if (this.containerProperties.getAckMode() == AckMode.RECORD) { + if (this.ackMode.equals(AckMode.RECORD)) { handleAck(message); } - if (inRetryMode.get()) { - inRetryMode.set(false); - } + inRetryMode.compareAndSet(true, false); } catch (Exception e) { if (this.pulsarConsumerErrorHandler != null) { invokeRecordListenerErrorHandler(inRetryMode, message, e); } else { - if (this.containerProperties.getAckMode() == AckMode.RECORD) { + if (this.ackMode.equals(AckMode.RECORD)) { this.consumer.negativeAcknowledge(message); } - else if (this.containerProperties.getAckMode() == AckMode.BATCH) { + else if (this.ackMode.equals(AckMode.BATCH)) { this.nackableMessages.add(message.getMessageId()); } else { @@ -371,7 +374,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess while (inRetryMode.get()); } // All the records are processed at this point. Handle acks. - if (this.containerProperties.getAckMode() == AckMode.BATCH) { + if (this.ackMode.equals(AckMode.BATCH)) { handleAcks(messages); } } @@ -426,9 +429,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess inRetryMode.set(true); } else { - if (inRetryMode.get()) { - inRetryMode.set(false); - } + inRetryMode.compareAndSet(true, false); // retries exhausted - recover the message this.pulsarConsumerErrorHandler.recoverMessage(this.consumer, pulsarMessage, pulsarBatchListenerFailedException); @@ -453,14 +454,12 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess inRetryMode.set(true); } else { - if (inRetryMode.get()) { - inRetryMode.set(false); - } + inRetryMode.compareAndSet(true, false); // retries exhausted - recover the message this.pulsarConsumerErrorHandler.recoverMessage(this.consumer, message, e); // retries exhausted - if record ackmode, acknowledge, otherwise normal // batch ack at the end - if (this.containerProperties.getAckMode() == AckMode.RECORD) { + if (this.ackMode.equals(AckMode.RECORD)) { handleAck(message); } } @@ -468,12 +467,8 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess private void pendingMessagesHandledSuccessfully(AtomicBoolean inRetryMode, AtomicBoolean messagesPendingInBatch) { - if (inRetryMode.get()) { - inRetryMode.set(false); - } - if (messagesPendingInBatch.get()) { - messagesPendingInBatch.set(false); - } + inRetryMode.compareAndSet(true, false); + messagesPendingInBatch.compareAndSet(true, false); this.pulsarConsumerErrorHandler.clearMessage(); } @@ -483,8 +478,8 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess } private boolean isSharedSubscriptionType() { - return this.containerProperties.getSubscriptionType() == SubscriptionType.Shared - || this.containerProperties.getSubscriptionType() == SubscriptionType.Key_Shared; + return this.subscriptionType.equals(SubscriptionType.Shared) + || this.subscriptionType.equals(SubscriptionType.Key_Shared); } private void handleAcks(Messages messages) {