From 09024b24504bb5a0465109814d75e33e999456e1 Mon Sep 17 00:00:00 2001 From: Zixuan Liu Date: Wed, 15 Apr 2026 17:22:29 +0800 Subject: [PATCH 1/3] [improve][client] Best-effort retry for acks on send failure when ackReceiptEnabled=false Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- .../apache/pulsar/client/impl/ClientCnx.java | 30 +++--- ...sistentAcknowledgmentsGroupingTracker.java | 88 +++++++++++++++--- .../AcknowledgementsGroupingTrackerTest.java | 93 ++++++++++++++++++- .../ClientCnxRequestTimeoutQueueTest.java | 1 + .../client/impl/ClientTestFixtures.java | 5 + 5 files changed, 187 insertions(+), 30 deletions(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java index a699638632286..0c8bf66ae2522 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java @@ -1080,29 +1080,25 @@ CompletableFuture sendRequestWithId(ByteBuf cmd, long requestI } private void sendRequestAndHandleTimeout(ByteBuf requestMessage, long requestId, - RequestType requestType, boolean flush, - TimedCompletableFuture future) { + RequestType requestType, boolean flush, + TimedCompletableFuture future) { pendingRequests.put(requestId, future); - if (flush) { - ctx.writeAndFlush(requestMessage).addListener(writeFuture -> { - if (!writeFuture.isSuccess()) { - if (pendingRequests.remove(requestId, future) && !future.isDone()) { - log.warn() - .attr("send", requestType.getDescription()) - .exceptionMessage(writeFuture.cause()) - .log("Failed to send to broker"); - future.completeExceptionally(writeFuture.cause()); - } + (flush ? ctx.writeAndFlush(requestMessage) : ctx.write(requestMessage)).addListener(writeFuture -> { + if (!writeFuture.isSuccess()) { + if (pendingRequests.remove(requestId, future) && !future.isDone()) { + log.warn() + .attr("send", requestType.getDescription()) + .exceptionMessage(writeFuture.cause()) + .log("Failed to send to broker"); + future.completeExceptionally(writeFuture.cause()); } - }); - } else { - ctx.write(requestMessage, ctx().voidPromise()); - } + } + }); requestTimeoutQueue.add(new RequestTime(requestId, requestType)); } private CompletableFuture sendRequestAndHandleTimeout(ByteBuf requestMessage, long requestId, - RequestType requestType, boolean flush) { + RequestType requestType, boolean flush) { TimedCompletableFuture future = new TimedCompletableFuture<>(); sendRequestAndHandleTimeout(requestMessage, requestId, requestType, flush, future); return future; diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java index 0cd013f20b235..fd0308906de32 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java @@ -23,9 +23,11 @@ import io.netty.buffer.ByteBuf; import io.netty.channel.EventLoopGroup; import io.netty.util.concurrent.FastThreadLocal; +import java.util.AbstractMap.SimpleImmutableEntry; import java.util.ArrayList; import java.util.BitSet; import java.util.Collections; +import java.util.HashMap; import java.util.HashSet; import java.util.List; import java.util.Map; @@ -279,6 +281,11 @@ private CompletableFuture doIndividualAckAsync(MessageIdAdv messageId) { return CompletableFuture.completedFuture(null); } + @VisibleForTesting + int getPendingIndividualAcksSize() { + return pendingIndividualAcks.size(); + } + private CompletableFuture doIndividualBatchAck(MessageIdAdv batchMessageId, Map properties) { if (acknowledgementGroupTimeMicros == 0 || (properties != null && !properties.isEmpty())) { @@ -405,7 +412,7 @@ private CompletableFuture doImmediateBatchIndexAck(MessageIdAdv msgId, int } CompletableFuture completableFuture = newMessageAckCommandAndWrite(cnx, consumer.consumerId, - msgId.getLedgerId(), msgId.getEntryId(), bitSet, ackType, properties, true, null, null); + msgId.getLedgerId(), msgId.getEntryId(), bitSet, ackType, properties, true, null, null, null); bitSet.recycle(); return completableFuture; } @@ -437,16 +444,19 @@ private void flushAsync(ClientCnx cnx) { if (lastCumulativeAckToFlush != null) { shouldFlush = true; final MessageIdAdv messageId = lastCumulativeAckToFlush.getMessageId(); + MessageIdImpl[] chunkMsgIds = this.consumer.unAckedChunkedMessageIdSequenceMap.remove(messageId); newMessageAckCommandAndWrite(cnx, consumer.consumerId, messageId.getLedgerId(), messageId.getEntryId(), lastCumulativeAckToFlush.getBitSetRecyclable(), AckType.Cumulative, Collections.emptyMap(), false, - (TimedCompletableFuture) this.currentCumulativeAckFuture, null); - this.consumer.unAckedChunkedMessageIdSequenceMap.remove(messageId); + (TimedCompletableFuture) this.currentCumulativeAckFuture, null, + () -> restoreCumulativeAck(lastCumulativeAckToFlush, chunkMsgIds)); } // Flush all individual acks List> entriesToAck = new ArrayList<>(pendingIndividualAcks.size() + pendingIndividualBatchIndexAcks.size()); + List individualAcksToFlush = new ArrayList<>(pendingIndividualAcks.size()); + Map chunkedMessageIdsToRestore = new HashMap<>(); if (!pendingIndividualAcks.isEmpty()) { if (Commands.peerSupportsMultiMessageAcknowledgment(cnx.getRemoteEndpointProtocolVersion())) { // We can send 1 single protobuf command with all individual acks @@ -455,11 +465,13 @@ private void flushAsync(ClientCnx cnx) { if (msgId == null) { break; } + individualAcksToFlush.add(msgId); // if messageId is checked then all the chunked related to that msg also processed so, ack all of // them MessageIdImpl[] chunkMsgIds = this.consumer.unAckedChunkedMessageIdSequenceMap.get(msgId); if (chunkMsgIds != null && chunkMsgIds.length > 1) { + chunkedMessageIdsToRestore.put(msgId, chunkMsgIds); for (MessageIdImpl cMsgId : chunkMsgIds) { if (cMsgId != null) { entriesToAck.add(Triple.of(cMsgId.getLedgerId(), cMsgId.getEntryId(), null)); @@ -478,14 +490,16 @@ private void flushAsync(ClientCnx cnx) { if (msgId == null) { break; } + individualAcksToFlush.add(msgId); newMessageAckCommandAndWrite(cnx, consumer.consumerId, msgId.getLedgerId(), msgId.getEntryId(), null, AckType.Individual, Collections.emptyMap(), false, - null, null); + null, null, () -> restoreIndividualAck(msgId, null)); shouldFlush = true; } } } + List> batchIndexAcksToFlush = new ArrayList<>(); while (true) { Map.Entry entry = pendingIndividualBatchIndexAcks.pollFirstEntry(); @@ -493,6 +507,7 @@ private void flushAsync(ClientCnx cnx) { // The entry has been removed in a different thread break; } + batchIndexAcksToFlush.add(entry); entriesToAck.add(Triple.of( entry.getKey().getLedgerId(), entry.getKey().getEntryId(), entry.getValue())); } @@ -501,7 +516,9 @@ private void flushAsync(ClientCnx cnx) { newMessageAckCommandAndWrite(cnx, consumer.consumerId, 0L, 0L, null, AckType.Individual, null, true, - (TimedCompletableFuture) currentIndividualAckFuture, entriesToAck); + (TimedCompletableFuture) currentIndividualAckFuture, entriesToAck, + () -> restoreIndividualAndBatchIndexAcks(individualAcksToFlush, chunkedMessageIdsToRestore, + batchIndexAcksToFlush)); shouldFlush = true; } @@ -546,19 +563,19 @@ private CompletableFuture newImmediateAckAndFlush(long consumerId, Message } } completableFuture = newMessageAckCommandAndWrite(cnx, consumer.consumerId, 0L, 0L, - null, ackType, null, true, null, entriesToAck); + null, ackType, null, true, null, entriesToAck, null); } else { // if don't support multi message ack, it also support ack receipt, so we should not think about the // ack receipt in this logic for (MessageIdImpl cMsgId : chunkMsgIds) { newMessageAckCommandAndWrite(cnx, consumerId, cMsgId.getLedgerId(), cMsgId.getEntryId(), - bitSet, ackType, map, true, null, null); + bitSet, ackType, map, true, null, null, null); } completableFuture = CompletableFuture.completedFuture(null); } } else { completableFuture = newMessageAckCommandAndWrite(cnx, consumerId, msgId.getLedgerId(), msgId.getEntryId(), - bitSet, ackType, map, true, null, null); + bitSet, ackType, map, true, null, null, null); } return completableFuture; } @@ -568,7 +585,8 @@ private CompletableFuture newMessageAckCommandAndWrite( long entryId, BitSetRecyclable ackSet, AckType ackType, Map properties, boolean flush, TimedCompletableFuture timedCompletableFuture, - List> entriesToAck) { + List> entriesToAck, + Runnable writeFailureCallback) { if (consumer.isAckReceiptEnabled()) { final long requestId = consumer.getClient().newRequestId(); final ByteBuf cmd; @@ -610,14 +628,62 @@ private CompletableFuture newMessageAckCommandAndWrite( cmd = Commands.newMultiMessageAck(consumerId, entriesToAck, -1); } if (flush) { - cnx.ctx().writeAndFlush(cmd, cnx.ctx().voidPromise()); + if (writeFailureCallback == null) { + cnx.ctx().writeAndFlush(cmd, cnx.ctx().voidPromise()); + } else { + cnx.ctx().writeAndFlush(cmd).addListener(writeFuture -> { + if (!writeFuture.isSuccess()) { + writeFailureCallback.run(); + } + }); + } } else { - cnx.ctx().write(cmd, cnx.ctx().voidPromise()); + if (writeFailureCallback == null) { + cnx.ctx().write(cmd, cnx.ctx().voidPromise()); + } else { + cnx.ctx().write(cmd).addListener(writeFuture -> { + if (!writeFuture.isSuccess()) { + writeFailureCallback.run(); + } + }); + } } return CompletableFuture.completedFuture(null); } } + private void restoreCumulativeAck(LastCumulativeAck ackToRestore, @Nullable MessageIdImpl[] chunkMsgIds) { + lastCumulativeAck.update(ackToRestore.getMessageId(), cloneBitSet(ackToRestore.getBitSetRecyclable())); + restoreChunkedMessageIds(ackToRestore.getMessageId(), chunkMsgIds); + } + + private void restoreIndividualAck(MessageIdAdv messageId, @Nullable MessageIdImpl[] chunkMsgIds) { + pendingIndividualAcks.add(messageId); + restoreChunkedMessageIds(messageId, chunkMsgIds); + } + + private void restoreIndividualAndBatchIndexAcks(List messageIds, + Map chunkMsgIds, + List> batchIndexAcks) { + pendingIndividualAcks.addAll(messageIds); + chunkMsgIds.forEach(this::restoreChunkedMessageIds); + batchIndexAcks.forEach(entry -> pendingIndividualBatchIndexAcks.merge(entry.getKey(), entry.getValue(), + (currentValue, valueToRestore) -> { + currentValue.and(valueToRestore); + return currentValue; + })); + } + + private void restoreChunkedMessageIds(MessageIdAdv messageId, @Nullable MessageIdImpl[] chunkMsgIds) { + if (chunkMsgIds != null) { + consumer.unAckedChunkedMessageIdSequenceMap.putIfAbsent(messageId, chunkMsgIds); + } + } + + private BitSetRecyclable cloneBitSet(@Nullable BitSetRecyclable bitSet) { + return bitSet != null ? BitSetRecyclable.valueOf(bitSet.toLongArray()) : null; + } + public Optional acquireReadLock() { Optional optionalLock = Optional.ofNullable(consumer.isAckReceiptEnabled() ? lock.readLock() : null); optionalLock.ifPresent(Lock::lock); diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/AcknowledgementsGroupingTrackerTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/AcknowledgementsGroupingTrackerTest.java index 0b91d86064cde..8bb85b054a3e3 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/AcknowledgementsGroupingTrackerTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/AcknowledgementsGroupingTrackerTest.java @@ -18,6 +18,7 @@ */ package org.apache.pulsar.client.impl; +import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.mock; @@ -28,9 +29,13 @@ import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertTrue; import io.netty.buffer.ByteBuf; +import io.netty.channel.ChannelFuture; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.EventLoopGroup; import io.netty.channel.nio.NioEventLoopGroup; +import io.netty.util.ReferenceCountUtil; +import io.netty.util.concurrent.Future; +import io.netty.util.concurrent.GenericFutureListener; import java.util.ArrayList; import java.util.BitSet; import java.util.Collections; @@ -56,6 +61,8 @@ public class AcknowledgementsGroupingTrackerTest { private ConsumerImpl consumer; private EventLoopGroup eventLoopGroup; private AtomicBoolean returnCnx = new AtomicBoolean(true); + private ChannelHandlerContext successCtx; + private AtomicBoolean failAckCommandSend = new AtomicBoolean(false); @BeforeClass public void setup() throws NoSuchFieldException, IllegalAccessException { @@ -70,9 +77,9 @@ public void setup() throws NoSuchFieldException, IllegalAccessException { doReturn(new ConsumerStatsRecorderImpl()).when(consumer).getStats(); doReturn(UnAckedMessageTracker.UNACKED_MESSAGE_TRACKER_DISABLED) .when(consumer).getUnAckedMessageTracker(); - ChannelHandlerContext ctx = ClientTestFixtures.mockChannelHandlerContext(); + successCtx = ClientTestFixtures.mockChannelHandlerContext(); doAnswer(invocation -> returnCnx.get() ? cnx : null).when(consumer).getClientCnx(); - doReturn(ctx).when(cnx).ctx(); + doReturn(successCtx).when(cnx).ctx(); } @DataProvider(name = "isNeedReceipt") @@ -323,6 +330,60 @@ public void testAckTrackerMultiAck(boolean isNeedReceipt) { tracker.close(); } + @Test + public void testFlushRetainsPendingIndividualAckOnSendFailureWithoutAckReceipt() throws Exception { + ConsumerConfigurationData conf = new ConsumerConfigurationData<>(); + conf.setAcknowledgementsGroupTimeMicros(TimeUnit.SECONDS.toMicros(10)); + conf.setAckReceiptEnabled(false); + doReturn(false).when(consumer).isAckReceiptEnabled(); + PersistentAcknowledgmentsGroupingTracker tracker = + new PersistentAcknowledgmentsGroupingTracker(consumer, conf, eventLoopGroup); + + MessageIdImpl msg1 = new MessageIdImpl(5, 1, 0); + tracker.addAcknowledgment(msg1, AckType.Individual, Collections.emptyMap()); + assertEquals(tracker.getPendingIndividualAcksSize(), 1); + + doReturn(createFailedChannelHandlerContext()).when(cnx).ctx(); + + tracker.flush(); + + assertTrue(tracker.isDuplicate(msg1)); + assertEquals(tracker.getPendingIndividualAcksSize(), 1); + + doReturn(successCtx).when(cnx).ctx(); + + tracker.flush(); + + assertFalse(tracker.isDuplicate(msg1)); + assertEquals(tracker.getPendingIndividualAcksSize(), 0); + tracker.close(); + } + + @Test + public void testFlushFailsAckFutureOnSendFailureWithAckReceipt() throws Exception { + ConsumerConfigurationData conf = new ConsumerConfigurationData<>(); + conf.setAcknowledgementsGroupTimeMicros(TimeUnit.SECONDS.toMicros(10)); + conf.setAckReceiptEnabled(true); + doReturn(true).when(consumer).isAckReceiptEnabled(); + PersistentAcknowledgmentsGroupingTracker tracker = + new PersistentAcknowledgmentsGroupingTracker(consumer, conf, eventLoopGroup); + + MessageIdImpl msg1 = new MessageIdImpl(5, 1, 0); + CompletableFuture ackFuture = + tracker.addAcknowledgment(msg1, AckType.Individual, Collections.emptyMap()); + assertEquals(tracker.getPendingIndividualAcksSize(), 1); + + failAckCommandSend.set(true); + tracker.flush(); + + assertTrue(ackFuture.isCompletedExceptionally()); + assertFalse(tracker.isDuplicate(msg1)); + assertEquals(tracker.getPendingIndividualAcksSize(), 0); + + failAckCommandSend.set(false); + tracker.close(); + } + @Test(dataProvider = "isNeedReceipt") public void testBatchAckTrackerMultiAck(boolean isNeedReceipt) throws Exception { ConsumerConfigurationData conf = new ConsumerConfigurationData<>(); @@ -463,12 +524,40 @@ public ClientCnxTest(ClientConfigurationData conf, EventLoopGroup eventLoopGroup @Override public CompletableFuture newAckForReceipt(ByteBuf request, long requestId) { + if (failAckCommandSend.get()) { + return CompletableFuture.failedFuture(new RuntimeException("ack send failed")); + } return CompletableFuture.completedFuture(null); } @Override public void newAckForReceiptWithFuture(ByteBuf request, long requestId, TimedCompletableFuture future) { + if (failAckCommandSend.get()) { + future.completeExceptionally(new RuntimeException("ack send failed")); + } } } + + private ChannelHandlerContext createFailedChannelHandlerContext() { + ChannelHandlerContext ctx = mock(ChannelHandlerContext.class); + ChannelFuture listenerFuture = mock(ChannelFuture.class); + ChannelFuture failedFuture = mock(ChannelFuture.class); + when(failedFuture.isSuccess()).thenReturn(false); + when(failedFuture.cause()).thenReturn(new RuntimeException("ack send failed")); + doAnswer(invocation -> { + GenericFutureListener> listener = invocation.getArgument(0); + listener.operationComplete(failedFuture); + return listenerFuture; + }).when(listenerFuture).addListener(any()); + doAnswer(invocation -> { + ReferenceCountUtil.release(invocation.getArgument(0)); + return listenerFuture; + }).when(ctx).write(any()); + doAnswer(invocation -> { + ReferenceCountUtil.release(invocation.getArgument(0)); + return listenerFuture; + }).when(ctx).writeAndFlush(any()); + return ctx; + } } diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ClientCnxRequestTimeoutQueueTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ClientCnxRequestTimeoutQueueTest.java index c4dab1ad351ab..7e5d98b136a40 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ClientCnxRequestTimeoutQueueTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ClientCnxRequestTimeoutQueueTest.java @@ -64,6 +64,7 @@ void setupClientCnx() throws Exception { ChannelHandlerContext ctx = mock(ChannelHandlerContext.class); Channel channel = mock(Channel.class); when(ctx.writeAndFlush(any())).thenAnswer(args -> mock(ChannelFuture.class)); + when(ctx.write(any())).thenAnswer(args -> mock(ChannelFuture.class)); when(ctx.channel()).thenReturn(channel); when(channel.remoteAddress()).thenReturn(new InetSocketAddress(1234)); cnx.channelActive(ctx); diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ClientTestFixtures.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ClientTestFixtures.java index e4ab5da96f649..c8bcfc65a955c 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ClientTestFixtures.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ClientTestFixtures.java @@ -151,6 +151,11 @@ public static ChannelHandlerContext mockChannelHandlerContext() { }).when(listenerFuture).addListener(any()); // handle write and writeAndFlush methods so that the input message is released + doAnswer(invocation -> { + Object msg = invocation.getArgument(0); + ReferenceCountUtil.release(msg); + return listenerFuture; + }).when(ctx).write(any()); doAnswer(invocation -> { Object msg = invocation.getArgument(0); ReferenceCountUtil.release(msg); From c22da98585726bf16ef3fd4b744717eeeff8e23f Mon Sep 17 00:00:00 2001 From: Zixuan Liu Date: Wed, 15 Apr 2026 17:29:58 +0800 Subject: [PATCH 2/3] [improve][client] Skip cumulative ack restore on send failure Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- .../PersistentAcknowledgmentsGroupingTracker.java | 14 ++------------ 1 file changed, 2 insertions(+), 12 deletions(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java index fd0308906de32..1d91230992bd1 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java @@ -444,12 +444,11 @@ private void flushAsync(ClientCnx cnx) { if (lastCumulativeAckToFlush != null) { shouldFlush = true; final MessageIdAdv messageId = lastCumulativeAckToFlush.getMessageId(); - MessageIdImpl[] chunkMsgIds = this.consumer.unAckedChunkedMessageIdSequenceMap.remove(messageId); newMessageAckCommandAndWrite(cnx, consumer.consumerId, messageId.getLedgerId(), messageId.getEntryId(), lastCumulativeAckToFlush.getBitSetRecyclable(), AckType.Cumulative, Collections.emptyMap(), false, - (TimedCompletableFuture) this.currentCumulativeAckFuture, null, - () -> restoreCumulativeAck(lastCumulativeAckToFlush, chunkMsgIds)); + (TimedCompletableFuture) this.currentCumulativeAckFuture, null, null); + this.consumer.unAckedChunkedMessageIdSequenceMap.remove(messageId); } // Flush all individual acks @@ -652,11 +651,6 @@ private CompletableFuture newMessageAckCommandAndWrite( } } - private void restoreCumulativeAck(LastCumulativeAck ackToRestore, @Nullable MessageIdImpl[] chunkMsgIds) { - lastCumulativeAck.update(ackToRestore.getMessageId(), cloneBitSet(ackToRestore.getBitSetRecyclable())); - restoreChunkedMessageIds(ackToRestore.getMessageId(), chunkMsgIds); - } - private void restoreIndividualAck(MessageIdAdv messageId, @Nullable MessageIdImpl[] chunkMsgIds) { pendingIndividualAcks.add(messageId); restoreChunkedMessageIds(messageId, chunkMsgIds); @@ -680,10 +674,6 @@ private void restoreChunkedMessageIds(MessageIdAdv messageId, @Nullable MessageI } } - private BitSetRecyclable cloneBitSet(@Nullable BitSetRecyclable bitSet) { - return bitSet != null ? BitSetRecyclable.valueOf(bitSet.toLongArray()) : null; - } - public Optional acquireReadLock() { Optional optionalLock = Optional.ofNullable(consumer.isAckReceiptEnabled() ? lock.readLock() : null); optionalLock.ifPresent(Lock::lock); From be6d53bcc84155d94999b327091d3fcd1f68d96b Mon Sep 17 00:00:00 2001 From: Zixuan Liu Date: Wed, 15 Apr 2026 17:49:57 +0800 Subject: [PATCH 3/3] Remove unused import --- .../client/impl/PersistentAcknowledgmentsGroupingTracker.java | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java index 1d91230992bd1..001137bd912a9 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java @@ -23,7 +23,6 @@ import io.netty.buffer.ByteBuf; import io.netty.channel.EventLoopGroup; import io.netty.util.concurrent.FastThreadLocal; -import java.util.AbstractMap.SimpleImmutableEntry; import java.util.ArrayList; import java.util.BitSet; import java.util.Collections;