From 01e96ddfa35833d5c69ea1fe9d12a95d592df942 Mon Sep 17 00:00:00 2001 From: Chris Bono Date: Fri, 19 Apr 2024 11:22:02 -0500 Subject: [PATCH] Remove duplicate handleBatchAcks method --- ...DefaultPulsarMessageListenerContainer.java | 59 ++++++++----------- 1 file changed, 25 insertions(+), 34 deletions(-) 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 fc690ac7..6ba4b59d 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 @@ -529,7 +529,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess } // All the records are processed at this point - handle acks if (this.ackMode.equals(AckMode.BATCH)) { - handleBatchAcks(messages, null); + handleBatchAcksForRecordListener(messages, null); } } @@ -670,21 +670,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess this.batchMessageListener.received(this.consumer, messageList); } if (this.ackMode.equals(AckMode.BATCH)) { - try { - if (isSharedSubscriptionType()) { - AckUtils.handleAck(this.consumer, messages, txn); - } - else { - Stream> stream = StreamSupport.stream(messages.spliterator(), true); - Message last = stream.reduce((a, b) -> b).orElse(null); - AckUtils.handleAckCumulative(this.consumer, last, txn); - } - } - catch (PulsarException pe) { - DefaultPulsarMessageListenerContainer.this.logger.warn(pe, - () -> "Batch acknowledgment failed: " + pe.getMessage()); - this.consumer.negativeAcknowledge(messages); - } + handleBatchAcks(messages, txn); } if (this.pulsarConsumerErrorHandler != null) { pendingMessagesHandledSuccessfully(inRetryMode, messagesPendingInBatch); @@ -793,25 +779,9 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess || this.subscriptionType.equals(SubscriptionType.Key_Shared)); } - private void handleBatchAcks(Messages messages, @Nullable Transaction txn) { + private void handleBatchAcksForRecordListener(Messages messages, @Nullable Transaction txn) { if (this.nackableMessages.isEmpty()) { - try { - if (messages.size() > 0) { - if (isSharedSubscriptionType()) { - AckUtils.handleAck(this.consumer, messages, txn); - } - else { - Stream> stream = StreamSupport.stream(messages.spliterator(), true); - Message last = stream.reduce((a, b) -> b).orElse(null); - AckUtils.handleAckCumulative(this.consumer, last, txn); - } - } - } - catch (PulsarException pe) { - DefaultPulsarMessageListenerContainer.this.logger.warn(pe, - () -> "Batch acks failed: " + pe.getMessage()); - this.consumer.negativeAcknowledge(messages); - } + handleBatchAcks(messages, txn); } else { for (Message message : messages) { @@ -826,6 +796,27 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess } } + private void handleBatchAcks(Messages messages, @Nullable Transaction txn) { + if (messages.size() <= 0) { + return; + } + try { + if (isSharedSubscriptionType()) { + AckUtils.handleAck(this.consumer, messages, txn); + } + else { + Stream> stream = StreamSupport.stream(messages.spliterator(), true); + Message last = stream.reduce((a, b) -> b).orElse(null); + AckUtils.handleAckCumulative(this.consumer, last, txn); + } + } + catch (PulsarException pe) { + DefaultPulsarMessageListenerContainer.this.logger.warn(pe, + () -> "Batch acknowledgment failed: " + pe.getMessage()); + this.consumer.negativeAcknowledge(messages); + } + } + private void handleAck(Message message, @Nullable Transaction txn) { AckUtils.handleAckWithNackOnFailure(this.consumer, message.getMessageId(), txn); }