Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -312,6 +312,104 @@ private void sendSingleChunk(Producer<String> 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<String> consumer = pulsarClient.newConsumer(Schema.STRING)
.topic(topicName)
.subscriptionName(subName)
.expireTimeOfIncompleteChunkedMessage(1, TimeUnit.SECONDS)
.subscribe();
@Cleanup
Producer<String> producer = pulsarClient.newProducer(Schema.STRING)
.topic(topicName)
.chunkMaxMessageSize(100)
.enableChunking(true)
.enableBatching(false)
.create();

ConsumerImpl<String> consumerImpl = (ConsumerImpl<String>) 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<String> 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"));
}

/**
* 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<String> consumer = pulsarClient.newConsumer(Schema.STRING)
.topic(topicName)
.subscriptionName("my-sub")
.expireTimeOfIncompleteChunkedMessage(1, TimeUnit.SECONDS)
.subscribe();
@Cleanup
Producer<String> producer = pulsarClient.newProducer(Schema.STRING)
.topic(topicName)
.chunkMaxMessageSize(100)
.enableChunking(true)
.enableBatching(false)
.create();

ConsumerImpl<String> consumerImpl = (ConsumerImpl<String>) 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
Expand Down Expand Up @@ -431,6 +529,56 @@ public void testOrphanedLastChunkDoesNotStallDispatch() throws Exception {
consumer.acknowledge(msg);
}

/**
* 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<String> consumer = pulsarClient.newConsumer(Schema.STRING)
.topic(topicName)
.subscriptionName(subName)
.maxPendingChunkedMessage(10)
.expireTimeOfIncompleteChunkedMessage(1, TimeUnit.HOURS)
.autoAckOldestChunkedMessageOnQueueFull(true)
.subscribe();
@Cleanup
Producer<String> producer = pulsarClient.newProducer(Schema.STRING)
.topic(topicName)
.chunkMaxMessageSize(100)
.enableChunking(true)
.enableBatching(false)
.create();

ConsumerImpl<String> consumerImpl = (ConsumerImpl<String>) 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<String> 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");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -228,6 +228,11 @@ public class ConsumerImpl<T> extends ConsumerBase<T> implements ConnectionHandle
// it will be used to manage N outstanding chunked message buffers
private final BlockingQueue<String> pendingChunkedMessageUuidQueue;

@VisibleForTesting
int getPendingChunkedMessageUuidQueueSizeForTest() {
return pendingChunkedMessageUuidQueue.size();
}

private final boolean createTopicIfDoesNotExist;
private final boolean poolMessages;

Expand Down Expand Up @@ -1535,6 +1540,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());
Comment thread
programmerahul marked this conversation as resolved.
chunkedMsgCtx.recycle();
}

Expand Down Expand Up @@ -3191,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;
}
}
Expand Down