This is an automated email from the ASF dual-hosted git repository.
lhotari pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pulsar.git
The following commit(s) were added to refs/heads/master by this push:
new e8df873ddfd [fix][client] Sync ackSet in client with broker to stop
acked messages reaching the DLQ (#26135)
e8df873ddfd is described below
commit e8df873ddfde4d8b5a38a49e4b8765c693d65ab9
Author: Yike Xiao <[email protected]>
AuthorDate: Fri Jul 3 21:34:38 2026 +0800
[fix][client] Sync ackSet in client with broker to stop acked messages
reaching the DLQ (#26135)
---
.../pulsar/client/api/DeadLetterTopicTest.java | 72 ++++++++++++++++++++++
.../apache/pulsar/client/impl/ConsumerImpl.java | 22 ++++---
.../client/impl/MessagePayloadContextImpl.java | 9 ++-
.../pulsar/client/impl/ZeroQueueConsumerImpl.java | 2 +-
.../pulsar/client/impl/ConsumerImplTest.java | 57 +++++++++++++++++
5 files changed, 149 insertions(+), 13 deletions(-)
diff --git
a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/DeadLetterTopicTest.java
b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/DeadLetterTopicTest.java
index 665def299ba..4d15498c2b0 100644
---
a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/DeadLetterTopicTest.java
+++
b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/DeadLetterTopicTest.java
@@ -1689,4 +1689,76 @@ public class DeadLetterTopicTest extends
SharedPulsarBaseTest {
verify(client, times(0)).getPartitionedTopicMetadata(anyString(),
anyBoolean(), anyBoolean());
}
+ @Test
+ public void
testAckedBatchMessageNotSentToDeadLetterTopicOnFinalRedeliveryRound() throws
Exception {
+ final String topic = newTopicName();
+ final int maxRedeliveryCount = 3;
+ final int batchSize = 5;
+ final String subscriptionName = "my-subscription";
+
+ Consumer<byte[]> consumer = pulsarClient.newConsumer(Schema.BYTES)
+ .topic(topic)
+ .subscriptionName(subscriptionName)
+ .subscriptionType(SubscriptionType.Shared)
+
.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest)
+ .enableBatchIndexAcknowledgment(true)
+
.deadLetterPolicy(DeadLetterPolicy.builder().maxRedeliverCount(maxRedeliveryCount).build())
+ .ackTimeout(1, TimeUnit.SECONDS)
+ .receiverQueueSize(100)
+ .subscribe();
+
+ @Cleanup
+ PulsarClient newPulsarClient = newPulsarClient();
+ Consumer<byte[]> deadLetterConsumer =
newPulsarClient.newConsumer(Schema.BYTES)
+ .topic(topic + "-" + subscriptionName + "-DLQ")
+ .subscriptionName(subscriptionName)
+
.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest)
+ .subscribe();
+
+ Producer<byte[]> producer = pulsarClient.newProducer(Schema.BYTES)
+ .topic(topic)
+ .enableBatching(true)
+ .batchingMaxPublishDelay(1, TimeUnit.SECONDS)
+ .create();
+ List<CompletableFuture<MessageId>> sendFutures = new ArrayList<>();
+ for (int i = 0; i < batchSize; i++) {
+ sendFutures.add(producer.newMessage().value(("message-" +
i).getBytes()).sendAsync());
+ }
+ for (CompletableFuture<MessageId> future : sendFutures) {
+ future.get();
+ }
+ producer.close();
+
+ // Batch indices 1 and 2 are deliberately left to time out for the
first `maxRedeliveryCount`
+ // rounds, then explicitly acked on the final round (redeliveryCount
== maxRedeliveryCount) --
+ // the app's last chance to prevent them from being routed to the DLQ.
The other 3 messages in
+ // the batch are acked immediately on the first delivery.
+ // Expected deliveries: (batchSize - 2) once each, plus indices 1 and
2 redelivered on every
+ // round from 0 through maxRedeliveryCount inclusive.
+ final int expectedDeliveries = (batchSize - 2) + 2 *
(maxRedeliveryCount + 1);
+ int received = 0;
+ while (received < expectedDeliveries) {
+ Message<byte[]> message = consumer.receive(5, TimeUnit.SECONDS);
+ assertNotNull(message, "consumer should keep receiving messages
until the batch settles");
+ received++;
+ MessageIdAdv messageId = (MessageIdAdv) message.getMessageId();
+ int batchIndex = messageId.getBatchIndex();
+ int redeliveryCount = message.getRedeliveryCount();
+ if ((batchIndex == 1 || batchIndex == 2) && redeliveryCount <
maxRedeliveryCount) {
+ // Let it time out instead of acking.
+ continue;
+ }
+ consumer.acknowledge(message);
+ }
+
+ // No message should ever be routed to the DLQ, since every message
was explicitly acked at or
+ // before its final allowed redelivery round.
+ Message<byte[]> deadLetterMessage = deadLetterConsumer.receive(5,
TimeUnit.SECONDS);
+ assertNull(deadLetterMessage, "no message should have been routed to
the DLQ, "
+ + "but received: " + deadLetterMessage);
+
+ deadLetterConsumer.close();
+ consumer.close();
+ }
+
}
diff --git
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java
index a133826842d..cf5f2a44527 100644
---
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java
+++
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java
@@ -23,6 +23,7 @@ import static
org.apache.pulsar.common.protocol.Commands.DEFAULT_CONSUMER_EPOCH;
import static org.apache.pulsar.common.protocol.Commands.hasChecksum;
import static org.apache.pulsar.common.protocol.Commands.serializeWithSize;
import static
org.apache.pulsar.common.util.Runnables.catchingAndLoggingThrowables;
+import static
org.apache.pulsar.common.util.SafeCollectionUtils.longArrayToList;
import com.google.common.annotations.VisibleForTesting;
import com.google.common.collect.ComparisonChain;
import com.google.common.collect.Iterables;
@@ -135,7 +136,6 @@ import org.apache.pulsar.common.util.Backoff;
import org.apache.pulsar.common.util.CompletableFutureCancellationHandler;
import org.apache.pulsar.common.util.ExceptionHandler;
import org.apache.pulsar.common.util.FutureUtil;
-import org.apache.pulsar.common.util.SafeCollectionUtils;
import org.apache.pulsar.common.util.collections.BitSetRecyclable;
import org.apache.pulsar.common.util.collections.ConcurrentBitSet;
import org.apache.pulsar.common.util.collections.GrowableArrayBlockingQueue;
@@ -144,6 +144,7 @@ public class ConsumerImpl<T> extends ConsumerBase<T>
implements ConnectionHandle
private static final Logger LOG = Logger.get(ConsumerImpl.class);
protected final Logger log;
+ private static final long[] EMPTY_ACK_SET = new long[0];
private static final int MAX_REDELIVER_UNACKNOWLEDGED = 1000;
final long consumerId;
@@ -1428,11 +1429,11 @@ public class ConsumerImpl<T> extends ConsumerBase<T>
implements ConnectionHandle
}
void messageReceived(CommandMessage cmdMessage, ByteBuf headersAndPayload,
ClientCnx cnx) {
- List<Long> ackSet = Collections.emptyList();
+ long[] ackSet = EMPTY_ACK_SET;
if (cmdMessage.getAckSetsCount() > 0) {
- ackSet = new ArrayList<>(cmdMessage.getAckSetsCount());
+ ackSet = new long[cmdMessage.getAckSetsCount()];
for (int i = 0; i < cmdMessage.getAckSetsCount(); i++) {
- ackSet.add(cmdMessage.getAckSetAt(i));
+ ackSet[i] = cmdMessage.getAckSetAt(i);
}
}
int redeliveryCount = cmdMessage.getRedeliveryCount();
@@ -1498,7 +1499,7 @@ public class ConsumerImpl<T> extends ConsumerBase<T>
implements ConnectionHandle
if (conf.getPayloadProcessor() != null) {
// uncompressedPayload is released in this method so we don't need
to call release() again
processPayloadByProcessor(brokerEntryMetadata, msgMetadata,
- uncompressedPayload, msgId, schema, redeliveryCount,
ackSet, consumerEpoch);
+ uncompressedPayload, msgId, schema, redeliveryCount,
longArrayToList(ackSet), consumerEpoch);
return;
}
@@ -1781,7 +1782,7 @@ public class ConsumerImpl<T> extends ConsumerBase<T>
implements ConnectionHandle
}
void receiveIndividualMessagesFromBatch(BrokerEntryMetadata
brokerEntryMetadata, MessageMetadata msgMetadata,
- int redeliveryCount, List<Long>
ackSet, ByteBuf uncompressedPayload,
+ int redeliveryCount, long[]
ackSet, ByteBuf uncompressedPayload,
MessageIdData messageId, ClientCnx
cnx, long consumerEpoch,
boolean isEncrypted) {
int batchSize = msgMetadata.getNumMessagesInBatch();
@@ -1796,8 +1797,9 @@ public class ConsumerImpl<T> extends ConsumerBase<T>
implements ConnectionHandle
BitSet ackSetInMessageId = BatchMessageIdImpl.newAckSet(batchSize);
BitSetRecyclable ackBitSet = null;
- if (ackSet != null && ackSet.size() > 0) {
- ackBitSet =
BitSetRecyclable.valueOf(SafeCollectionUtils.longListToArray(ackSet));
+ if (ackSet != null && ackSet.length > 0) {
+ ackBitSet = BitSetRecyclable.valueOf(ackSet);
+ ackSetInMessageId.and(BitSet.valueOf(ackSet));
}
SingleMessageMetadata singleMessageMetadata = new
SingleMessageMetadata();
@@ -2704,7 +2706,7 @@ public class ConsumerImpl<T> extends ConsumerBase<T>
implements ConnectionHandle
final ByteBuf seek;
if (msgId.getFirstChunkMessageId() != null) {
seek = Commands.newSeek(consumerId, requestId,
firstChunkMsgId.getLedgerId(),
- firstChunkMsgId.getEntryId(), new long[0]);
+ firstChunkMsgId.getEntryId(), EMPTY_ACK_SET);
} else {
final long[] ackSetArr;
if (MessageIdAdvUtils.isBatch(msgId)) {
@@ -2714,7 +2716,7 @@ public class ConsumerImpl<T> extends ConsumerBase<T>
implements ConnectionHandle
ackSetArr = ackSet.toLongArray();
ackSet.recycle();
} else {
- ackSetArr = new long[0];
+ ackSetArr = EMPTY_ACK_SET;
}
seek = Commands.newSeek(consumerId, requestId,
msgId.getLedgerId(), msgId.getEntryId(), ackSetArr);
}
diff --git
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagePayloadContextImpl.java
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagePayloadContextImpl.java
index 954afb57131..0220f1154d1 100644
---
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagePayloadContextImpl.java
+++
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagePayloadContextImpl.java
@@ -75,8 +75,13 @@ public class MessagePayloadContextImpl implements
MessagePayloadContext {
context.consumer = consumer;
context.redeliveryCount = redeliveryCount;
context.ackSetInMessageId =
BatchMessageIdImpl.newAckSet(context.getNumMessages());
- context.ackBitSet = (ackSet != null && ackSet.size() > 0)
- ?
BitSetRecyclable.valueOf(SafeCollectionUtils.longListToArray(ackSet))
+ boolean isAckSetNotEmpty = ackSet != null && ackSet.size() > 0;
+ long[] ackSetArray = SafeCollectionUtils.longListToArray(ackSet);
+ if (isAckSetNotEmpty) {
+ context.ackSetInMessageId.and(BitSet.valueOf(ackSetArray));
+ }
+ context.ackBitSet = isAckSetNotEmpty
+ ? BitSetRecyclable.valueOf(ackSetArray)
: null;
return context;
}
diff --git
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ZeroQueueConsumerImpl.java
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ZeroQueueConsumerImpl.java
index 5225bbd5fa2..dbbeed2973b 100644
---
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ZeroQueueConsumerImpl.java
+++
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ZeroQueueConsumerImpl.java
@@ -197,7 +197,7 @@ public class ZeroQueueConsumerImpl<T> extends
ConsumerImpl<T> {
@Override
void receiveIndividualMessagesFromBatch(BrokerEntryMetadata
brokerEntryMetadata, MessageMetadata msgMetadata,
- int redeliveryCount, List<Long>
ackSet, ByteBuf uncompressedPayload,
+ int redeliveryCount, long[]
ackSet, ByteBuf uncompressedPayload,
MessageIdData messageId, ClientCnx
cnx, long consumerEpoch,
boolean isEncrypted) {
diff --git
a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java
b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java
index 9d8b59db910..33732e56a5a 100644
---
a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java
+++
b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java
@@ -18,6 +18,7 @@
*/
package org.apache.pulsar.client.impl;
+import static
org.apache.pulsar.common.protocol.Commands.DEFAULT_CONSUMER_EPOCH;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.Mockito.any;
@@ -31,6 +32,10 @@ import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static org.testng.Assert.assertTrue;
import io.netty.buffer.ByteBuf;
+import io.netty.buffer.Unpooled;
+import java.util.Arrays;
+import java.util.BitSet;
+import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
import java.util.concurrent.ExecutorService;
@@ -38,17 +43,22 @@ import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.regex.Pattern;
+import java.util.stream.Collectors;
import lombok.Cleanup;
import org.apache.pulsar.client.api.Consumer;
import org.apache.pulsar.client.api.Message;
+import org.apache.pulsar.client.api.MessageIdAdv;
+import org.apache.pulsar.client.api.MessagePayload;
import org.apache.pulsar.client.api.Messages;
import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.client.api.PulsarClientException;
+import org.apache.pulsar.client.api.Schema;
import org.apache.pulsar.client.impl.conf.ClientConfigurationData;
import org.apache.pulsar.client.impl.conf.ConsumerConfigurationData;
import org.apache.pulsar.client.impl.conf.TopicConsumerConfigurationData;
import org.apache.pulsar.client.util.ExecutorProvider;
import org.apache.pulsar.client.util.ScheduledExecutorProvider;
+import org.apache.pulsar.common.api.proto.MessageMetadata;
import org.apache.pulsar.common.util.Backoff;
import org.awaitility.Awaitility;
import org.testng.Assert;
@@ -341,4 +351,51 @@ public class ConsumerImplTest {
"Hint must reflect the post-enqueue state (pipeline had >=1
message); "
+ "a concurrent drain of the just-enqueued message
must not clear it.");
}
+
+ @Test(invocationTimeOut = 1000)
+ public void testGetMessageAtSyncsAckSetInMessageIdWithBrokerAckSet() {
+ // Regression test for MessagePayloadContextImpl#getMessageAt: the
BatchMessageIdImpl handed
+ // back to the caller carries a shared ackSetInMessageId bitset that
must be seeded from the
+ // broker-reported ackSet, not a fresh "all unacked" bitset. Otherwise
indices the broker
+ // already knows are acked would be reported as still-outstanding in
the returned MessageId,
+ // which is the same root cause that let acked batch messages leak
into the DLQ (see
+ // ConsumerImpl#receiveIndividualMessagesFromBatch and its
ackSetInMessageId.and(...) fix).
+ final int batchSize = 3;
+ MessageMetadata messageMetadata = new MessageMetadata()
+ .setProducerName("test-producer")
+ .setSequenceId(0)
+ .setPublishTime(System.currentTimeMillis())
+ .setNumMessagesInBatch(batchSize);
+
+ // Broker reports index 0 as already acked (bit cleared); indices 1
and 2 are still
+ // outstanding (bits set). This mirrors the ackSet the broker attaches
on redelivery.
+ BitSet brokerAckSet = new BitSet(batchSize);
+ brokerAckSet.set(1);
+ brokerAckSet.set(2);
+ List<Long> ackSet =
Arrays.stream(brokerAckSet.toLongArray()).boxed().collect(Collectors.toList());
+
+ MessageIdImpl messageId = new MessageIdImpl(1L, 2L, -1);
+ MessagePayloadContextImpl context = MessagePayloadContextImpl.get(
+ null, messageMetadata, messageId, consumer, 0, ackSet,
DEFAULT_CONSUMER_EPOCH);
+ MessagePayload payload0 =
MessagePayloadImpl.create(Unpooled.wrappedBuffer(new byte[]{0}));
+ MessagePayload payload1 =
MessagePayloadImpl.create(Unpooled.wrappedBuffer(new byte[]{1}));
+ try {
+ // Index 0 is already acked per the broker, so it must not be
redelivered to the app.
+ Assert.assertNull(context.getMessageAt(0, batchSize, payload0,
false, Schema.BYTES));
+
+ Message<byte[]> message1 = context.getMessageAt(1, batchSize,
payload1, false, Schema.BYTES);
+ Assert.assertNotNull(message1);
+
+ BitSet ackSetInMessageId = ((MessageIdAdv)
message1.getMessageId()).getAckSet();
+ Assert.assertFalse(ackSetInMessageId.get(0),
+ "index 0 was already acked by the broker, so the returned
MessageId's ackSet "
+ + "must reflect it as acked, not fall back to the
default all-unacked state");
+ Assert.assertTrue(ackSetInMessageId.get(1), "index 1 is still
outstanding");
+ Assert.assertTrue(ackSetInMessageId.get(2), "index 2 is still
outstanding");
+ } finally {
+ payload0.release();
+ payload1.release();
+ context.recycle();
+ }
+ }
}