From 602cf88b4f621805e73aed45c221fb93a5ed6112 Mon Sep 17 00:00:00 2001 From: Rahul Prasad Date: Sun, 20 Sep 2026 00:24:41 +0530 Subject: [PATCH 1/3] uuid-queue --- .../client/impl/MessageChunkingTest.java | 50 +++++++++++++++++++ .../pulsar/client/impl/ConsumerImpl.java | 11 ++++ 2 files changed, 61 insertions(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java index 9616614d421b6..5f4a3222caffc 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java @@ -365,6 +365,56 @@ public void testMaxPendingChunkMessages() throws Exception { assertNull(consumer.receive(5, TimeUnit.SECONDS)); } + /** + * Verifies that pendingChunkedMessageUuidQueue does not leak entries as chunked messages + * complete normally. Each first chunk adds the message uuid to the queue; on completion the + * uuid must be removed so the queue stays in sync with chunkedMessagesMap. Without the fix the + * queue grew by one entry per completed chunked message unboundedly (a memory leak), since the + * only other removal paths (eviction/expiry) never run when maxPendingChunkedMessage is not + * exceeded. + */ + @Test + public void testPendingChunkedMessageUuidQueueDoesNotLeak() throws Exception { + final String topicName = "persistent://my-property/my-ns/uuidQueueNoLeak"; + final String subName = "my-sub"; + @Cleanup + Consumer consumer = pulsarClient.newConsumer(Schema.STRING) + .topic(topicName) + .subscriptionName(subName) + .maxPendingChunkedMessage(10) + .expireTimeOfIncompleteChunkedMessage(1, TimeUnit.HOURS) + .autoAckOldestChunkedMessageOnQueueFull(true) + .subscribe(); + @Cleanup + Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topicName) + .chunkMaxMessageSize(100) + .enableChunking(true) + .enableBatching(false) + .create(); + + ConsumerImpl consumerImpl = (ConsumerImpl) consumer; + + final int numMessages = 50; + for (int i = 0; i < numMessages; i++) { + String uuid = String.valueOf(i); + // A complete 2-chunk message. + sendSingleChunk(producer, uuid, 0, 2); + sendSingleChunk(producer, uuid, 1, 2); + Message msg = consumer.receive(5, TimeUnit.SECONDS); + assertEquals(msg.getValue(), "chunk-" + uuid + "-0|chunk-" + uuid + "-1|"); + consumer.acknowledge(msg); + } + + // Every message completed and was removed from chunkedMessagesMap; the uuid queue must have + // been drained in lockstep and not accumulated one ghost entry per completed message. + Awaitility.await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> { + assertEquals(consumerImpl.chunkedMessagesMap.size(), 0); + assertEquals(consumerImpl.getPendingChunkedMessageUuidQueueSizeForTest(), 0, + "pendingChunkedMessageUuidQueue leaked entries for completed chunked messages"); + }); + } + @Test public void testResendChunkMessagesWithoutAckHole() throws Exception { log.info().attr("method", methodName).log("Starting test"); 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 5862fce64de0e..072898b8abdb5 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 @@ -227,6 +227,11 @@ public class ConsumerImpl extends ConsumerBase implements ConnectionHandle // it will be used to manage N outstanding chunked message buffers private final BlockingQueue pendingChunkedMessageUuidQueue; + @VisibleForTesting + int getPendingChunkedMessageUuidQueueSizeForTest() { + return pendingChunkedMessageUuidQueue.size(); + } + private final boolean createTopicIfDoesNotExist; private final boolean poolMessages; @@ -1524,6 +1529,12 @@ void messageReceived(CommandMessage cmdMessage, ByteBuf headersAndPayload, Clien // add chunked messageId to unack-message tracker, and reduce pending-chunked-message count unAckedChunkedMessageIdSequenceMap.put(msgId, chunkedMsgCtx.chunkedMessageIds); pendingChunkedMessageCount--; + // The completed message's uuid was added to pendingChunkedMessageUuidQueue when its + // first chunk arrived, but is only ever removed by the eviction/expiry paths. On the + // normal completion path it was never removed, so the queue accumulated one entry per + // completed chunked message unboundedly (a memory leak for long-running consumers). + // Remove it here to keep the queue in sync with chunkedMessagesMap. + pendingChunkedMessageUuidQueue.remove(msgMetadata.getUuid()); chunkedMsgCtx.recycle(); } From f216708362117947297161af2eb7cab5e7e8b7f0 Mon Sep 17 00:00:00 2001 From: Rahul Prasad Date: Thu, 1 Oct 2026 18:52:16 +0530 Subject: [PATCH 2/3] added a testcase for expiry message mechanism being stuck --- .../client/impl/MessageChunkingTest.java | 58 +++++++++++++++++++ 1 file changed, 58 insertions(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java index 10185fee383b9..9a8c89fc943a4 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java @@ -312,6 +312,64 @@ private void sendSingleChunk(Producer producer, String uuid, int chunkId msg.send(); } + /** + * Verifies that expiry of incomplete chunked messages still works after an earlier chunked + * message has completed. + * + * removeExpireIncompleteChunkedMessages() only peek()s the head of pendingChunkedMessageUuidQueue + * and bails out (else -> return) the moment the head uuid is no longer present in + * chunkedMessagesMap. On the original client the completion path never removed a completed uuid + * from the queue, so the first completed message left a permanent "ghost" uuid at the head. From + * then on every expiry run saw that ghost, took the else branch and returned without ever + * inspecting the genuinely-incomplete uuids behind it: expiry was dead for the rest of the + * consumer's life (an unbounded leak of incomplete contexts). + * + * Here we complete one chunked message, then leave a second one incomplete. With the queue kept + * in sync on completion the incomplete message is expired and chunkedMessagesMap drains to empty; + * without it, the completed ghost blocks expiry and the incomplete context is never removed. + */ + @Test + public void testExpiryWorksAfterAnEarlierChunkedMessageCompletes() throws Exception { + final String topicName = "persistent://my-property/my-ns/expiryAfterComplete"; + final String subName = "my-sub"; + @Cleanup + Consumer consumer = pulsarClient.newConsumer(Schema.STRING) + .topic(topicName) + .subscriptionName(subName) + .expireTimeOfIncompleteChunkedMessage(1, TimeUnit.SECONDS) + .subscribe(); + @Cleanup + Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topicName) + .chunkMaxMessageSize(100) + .enableChunking(true) + .enableBatching(false) + .create(); + + ConsumerImpl consumerImpl = (ConsumerImpl) consumer; + + // 1. A complete chunked message. Receiving + acking it leaves its uuid as the head of + // pendingChunkedMessageUuidQueue unless completion keeps the queue in sync. + sendSingleChunk(producer, "done", 0, 2); + sendSingleChunk(producer, "done", 1, 2); + Message done = consumer.receive(5, TimeUnit.SECONDS); + assertNotNull(done); + assertEquals(done.getValue(), "chunk-done-0|chunk-done-1|"); + consumer.acknowledge(done); + + // 2. A second message left incomplete (only its first chunk arrives). Wait for the first + // chunk to be received and its assembly context to be registered before asserting expiry. + sendSingleChunk(producer, "stuck", 0, 2); + Awaitility.await().atMost(10, TimeUnit.SECONDS) + .untilAsserted(() -> assertEquals(consumerImpl.chunkedMessagesMap.size(), 1)); + + // 3. Past the expiry window, the incomplete "stuck" context must be discarded. If the + // completed "done" uuid still sits at the queue head, expiry bails and this never happens. + Awaitility.await().atMost(10, TimeUnit.SECONDS) + .untilAsserted(() -> assertEquals(consumerImpl.chunkedMessagesMap.size(), 0, + "expiry did not run: a completed message's uuid is blocking the queue head")); + } + /** * This test used to test the consumer configuration of maxPendingChunkedMessage. * If we set maxPendingChunkedMessage is 1 that means only one incomplete chunk message can be store in this From 27e45e5fa1353b14980fe8162927332af91a2962 Mon Sep 17 00:00:00 2001 From: Rahul Prasad Date: Fri, 2 Oct 2026 15:34:12 +0530 Subject: [PATCH 3/3] updated expire loop and added testcase for that --- .../client/impl/MessageChunkingTest.java | 40 +++++++++++++++++++ .../pulsar/client/impl/ConsumerImpl.java | 15 +++++-- 2 files changed, 52 insertions(+), 3 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java index 9a8c89fc943a4..872e16e153cf2 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java @@ -370,6 +370,46 @@ public void testExpiryWorksAfterAnEarlierChunkedMessageCompletes() throws Except "expiry did not run: a completed message's uuid is blocking the queue head")); } + /** + * Verifies that expiry still works when the queue head is a stale uuid left behind + * + * removeExpireIncompleteChunkedMessages() must poll past that ghost head rather than returning on + * it, otherwise an incomplete message queued behind the ghost would never expire. + */ + @Test + public void testExpirySkipsStaleQueueHeadFromDiscardPath() throws Exception { + final String topicName = "persistent://my-property/my-ns/expirySkipsDiscardGhost"; + @Cleanup + Consumer consumer = pulsarClient.newConsumer(Schema.STRING) + .topic(topicName) + .subscriptionName("my-sub") + .expireTimeOfIncompleteChunkedMessage(1, TimeUnit.SECONDS) + .subscribe(); + @Cleanup + Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topicName) + .chunkMaxMessageSize(100) + .enableChunking(true) + .enableBatching(false) + .create(); + + ConsumerImpl consumerImpl = (ConsumerImpl) consumer; + + // "ghost": chunk 0 arrives (queues the uuid), then chunk 2 of 3 arrives non-contiguously. + sendSingleChunk(producer, "ghost", 0, 3); + sendSingleChunk(producer, "ghost", 2, 3); + // "stuck": a genuinely incomplete message queued behind the ghost head. + sendSingleChunk(producer, "stuck", 0, 2); + Awaitility.await().atMost(10, TimeUnit.SECONDS) + .untilAsserted(() -> assertNotNull(consumerImpl.chunkedMessagesMap.get("stuck"))); + + // Past the expiry window, "stuck" must be collected. If the expiry loop returns on the ghost + // head instead of polling past it, "stuck" is never reached and stays in the map forever. + Awaitility.await().atMost(10, TimeUnit.SECONDS) + .untilAsserted(() -> assertNull(consumerImpl.chunkedMessagesMap.get("stuck"), + "expiry did not run: a stale discard-path uuid is blocking the queue head")); + } + /** * This test used to test the consumer configuration of maxPendingChunkedMessage. * If we set maxPendingChunkedMessage is 1 that means only one incomplete chunk message can be store in this 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 3adc3221a88af..5bf63e7cc1098 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 @@ -3202,15 +3202,24 @@ protected void removeExpireIncompleteChunkedMessages() { if (expireTimeOfIncompleteChunkedMessageMillis <= 0) { return; } - ChunkedMessageCtx chunkedMsgCtx = null; + ChunkedMessageCtx chunkedMsgCtx; String messageUUID; while ((messageUUID = pendingChunkedMessageUuidQueue.peek()) != null) { chunkedMsgCtx = StringUtils.isNotBlank(messageUUID) ? chunkedMessagesMap.get(messageUUID) : null; - if (chunkedMsgCtx != null && System - .currentTimeMillis() > (chunkedMsgCtx.receivedTime + expireTimeOfIncompleteChunkedMessageMillis)) { + if (chunkedMsgCtx == null) { + // Stale queue head with no backing context (e.g. the message already completed, or a + // discard path left the uuid behind). Like removeOldestPendingChunkedMessage(), skip + // it by polling and continue, so a ghost head cannot block expiry of live entries + // queued behind it. + pendingChunkedMessageUuidQueue.poll(); + continue; + } + if (System.currentTimeMillis() > (chunkedMsgCtx.receivedTime + expireTimeOfIncompleteChunkedMessageMillis)) { pendingChunkedMessageUuidQueue.remove(messageUUID); removeChunkMessage(messageUUID, chunkedMsgCtx, true); } else { + // Head is a live, not-yet-expired context. Entries are queued in arrival order, so + // nothing behind it can be older; stop here. return; } }