Repository navigation
[fix][client]in chunk message consumer, pendingChunkedMessageCount is not synced with count of entries present in chunkedMessagesMap #26689
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: master
Are you sure you want to change the base?
Changes from all commits
c81ea5d
0281fbf
ed03dcb
d4fa8db
54f6ab5
4d0711b
7b761e6
24d2088
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -439,6 +439,126 @@ public void testResendChunkMessages() throws Exception { | |
| Assert.assertEquals(((ConsumerImpl<String>) consumer).getAvailablePermits(), 8); | ||
| } | ||
|
|
||
| /** | ||
| * Proves that pendingChunkedMessageCount stays in sync with chunkedMessagesMap when duplicate | ||
| * (resent) first chunks arrive. | ||
| * | ||
| * When a first chunk (chunkId 0) arrives for a uuid that already has a ChunkedMessageCtx, | ||
| * processMessageChunk() removes the old context from chunkedMessagesMap without decrementing | ||
| * pendingChunkedMessageCount, then unconditionally increments it for the replacement. Before the | ||
| * fix, each duplicate first chunk inflated the count by 1 while the map size stayed 1 -- an | ||
| * upward drift that eventually pushes the count past maxPendingChunkedMessage and triggers | ||
| * spurious eviction of good in-flight messages. | ||
| * | ||
| * Invariant: pendingChunkedMessageCount == chunkedMessagesMap.size(). | ||
| * Buggy client: count == N, map == 1 -> FAIL. Fixed client: count == 1 -> PASS. | ||
| */ | ||
| @Test | ||
| public void testPendingChunkedMessageCountDriftOnDuplicateFirstChunk() throws Exception { | ||
| final String topicName = "persistent://my-property/my-ns/chunkCountDrift"; | ||
| @Cleanup | ||
| Consumer<String> consumer = pulsarClient.newConsumer(Schema.STRING) | ||
| .topic(topicName) | ||
| .subscriptionName("my-sub") | ||
| .maxPendingChunkedMessage(100) // high so eviction does not mask the drift | ||
| .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; | ||
|
|
||
| // Send the first chunk (chunkId 0) of the SAME uuid many times (duplicate/resent first | ||
| // chunk), never completing it. Each duplicate replaces the context in the map (size stays 1) | ||
| // but bumps pendingChunkedMessageCount on the buggy client. | ||
| final int duplicates = 10; | ||
| for (int i = 0; i < duplicates; i++) { | ||
| sendSingleChunk(producer, "dup-uuid", 0, 2); | ||
| } | ||
|
|
||
| Awaitility.await().atMost(10, TimeUnit.SECONDS) | ||
| .untilAsserted(() -> assertEquals(consumerImpl.chunkedMessagesMap.size(), 1)); | ||
|
|
||
| int mapSize = consumerImpl.chunkedMessagesMap.size(); | ||
| int count = consumerImpl.getPendingChunkedMessageCountForTest(); | ||
| int queueSize = consumerImpl.getPendingChunkedMessageUuidQueueSizeForTest(); | ||
|
|
||
| // All three structures track the same single in-flight uuid; they must stay consistent | ||
| // regardless of how many duplicate first chunks arrive for it. | ||
| assertEquals(queueSize, mapSize, | ||
| "pendingChunkedMessageUuidQueue.size (" + queueSize + ") drifted from chunkedMessagesMap.size (" | ||
| + mapSize + ") after " + duplicates + " duplicate first chunks"); | ||
| assertEquals(count, mapSize, | ||
| "pendingChunkedMessageCount (" + count + ") drifted from chunkedMessagesMap.size (" | ||
| + mapSize + ") after " + duplicates + " duplicate first chunks"); | ||
| } | ||
|
|
||
| /** | ||
| * A duplicate first chunk resets its context's receivedTime. If the uuid is left at its original | ||
| * position in pendingChunkedMessageUuidQueue, queue order no longer matches expiry order, and | ||
| * because removeExpireIncompleteChunkedMessages() only inspects the head and returns at the first | ||
| * non-expired entry, a repeatedly-refreshed head uuid blocks expiry of genuinely-expired entries | ||
| * behind it (the reviewer's A/B head-of-line scenario). | ||
| * | ||
| * Here A's first chunk is refreshed continuously so A never expires, while B is left untouched | ||
| * past the expiry window. B must still be expired and removed. The fix re-positions A's uuid in | ||
| * the queue on each refresh so B reaches the head and is collected. | ||
| */ | ||
| @Test | ||
| public void testRefreshedFirstChunkDoesNotBlockExpiryOfLaterEntries() throws Exception { | ||
| final String topicName = "persistent://my-property/my-ns/refreshedHeadExpiry"; | ||
| @Cleanup | ||
| Consumer<String> consumer = pulsarClient.newConsumer(Schema.STRING) | ||
| .topic(topicName) | ||
| .subscriptionName("my-sub") | ||
| .maxPendingChunkedMessage(100) | ||
| .expireTimeOfIncompleteChunkedMessage(2, 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; | ||
|
|
||
| // A arrives first, then B. Both incomplete. Queue: [A, B]. | ||
| sendSingleChunk(producer, "A", 0, 2); | ||
| sendSingleChunk(producer, "B", 0, 2); | ||
| Awaitility.await().atMost(10, TimeUnit.SECONDS) | ||
| .untilAsserted(() -> assertEquals(consumerImpl.chunkedMessagesMap.size(), 2)); | ||
|
|
||
| // Keep refreshing A's first chunk continuously so A's receivedTime is never older than the | ||
| // 2s expiry window -- A must never be eligible for expiry while we observe. We check B's | ||
| // removal WHILE still refreshing A: where the | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. [NIT] tidy the new expiry test's comments This comment has a fragment ( |
||
| // perpetually-refreshed head uuid actively shields the entry behind it. B is never refreshed, | ||
| // so once >2s pass it is expired -- unless A (left at the stale queue head without the fix) | ||
| // blocks the head-only expiry scan. Refresh every 400ms for up to ~12s, asserting B is gone. | ||
| long bArrivalNanos = System.nanoTime(); | ||
| boolean bCollected = false; | ||
| for (int i = 0; i < 30 && !bCollected; i++) { | ||
| sendSingleChunk(producer, "A", 0, 2); // keep A's receivedTime fresh (head stays non-expired) | ||
| Thread.sleep(400); | ||
| // only meaningful once B is well past its 2s expiry window | ||
| if (System.nanoTime() - bArrivalNanos > TimeUnit.SECONDS.toNanos(4)) { | ||
| bCollected = consumerImpl.chunkedMessagesMap.get("B") == null; | ||
| } | ||
| } | ||
|
|
||
| // B must have been expired and removed while A was still being refreshed at the (old) queue | ||
| // head. On the buggy client B stays stuck behind the refreshed A indefinitely. | ||
| assertNull(consumerImpl.chunkedMessagesMap.get("B"), | ||
| "expired message B was not collected while A was continually refreshed: a refreshed " | ||
| + "head uuid blocked head-only expiry of later entries"); | ||
| } | ||
|
|
||
| @Test | ||
| public void testExpireIncompleteChunkMessage() throws Exception{ | ||
| final String topicName = "persistent://my-property/my-ns/expireMsg"; | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -219,6 +219,12 @@ public class ConsumerImpl<T> extends ConsumerBase<T> implements ConnectionHandle | |
|
|
||
| protected Map<String, ChunkedMessageCtx> chunkedMessagesMap = new ConcurrentHashMap<>(); | ||
| private int pendingChunkedMessageCount = 0; | ||
|
|
||
| @VisibleForTesting | ||
| int getPendingChunkedMessageCountForTest() { | ||
| return pendingChunkedMessageCount; | ||
| } | ||
|
|
||
| protected long expireTimeOfIncompleteChunkedMessageMillis = 0; | ||
| private final AtomicBoolean expireChunkMessageTaskScheduled = new AtomicBoolean(false); | ||
| private final int maxPendingChunkedMessage; | ||
|
|
@@ -228,6 +234,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; | ||
|
|
||
|
|
@@ -1628,25 +1639,43 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m | |
| } | ||
| }); | ||
| } | ||
| // The first chunk of a new chunked-message received before receiving other chunks of previous | ||
| // chunked-message | ||
| // so, remove previous chunked-message from map and release buffer | ||
| // The first chunk of a new chunked-message received before receiving the other chunks | ||
| // of the previous chunked-message with the SAME uuid (a resent/duplicated first | ||
| // chunk, or a producer restart reusing the sequenceId). Discard the previous context. | ||
| // This is a REPLACEMENT of an existing tracking slot: the uuid is already counted once | ||
| // in pendingChunkedMessageCount, so this branch must NOT increment the count -- only | ||
| // the new-uuid branch below does. Otherwise each duplicate first chunk would drift the | ||
| // count above the real chunkedMessagesMap size (triggering spurious eviction). | ||
| // | ||
| // The uuid must, however, be re-positioned in pendingChunkedMessageUuidQueue. The | ||
| // replacement context below gets a fresh receivedTime, so leaving the uuid at its | ||
| // original (older) position would make queue order no longer match expiry order: | ||
| // removeExpireIncompleteChunkedMessages() only inspects the head and returns at the | ||
| // first non-expired entry, so a repeatedly-refreshed head uuid would indefinitely | ||
| // block expiry of genuinely-expired entries behind it. Remove the stale entry and | ||
| // re-add it so its position reflects the refreshed receivedTime. The count is | ||
| // unchanged (one entry out, one back in). | ||
| if (chunkedMsgCtx.chunkedMsgBuffer != null) { | ||
| ReferenceCountUtil.safeRelease(chunkedMsgCtx.chunkedMsgBuffer); | ||
| } | ||
| chunkedMsgCtx.recycle(); | ||
| chunkedMessagesMap.remove(msgMetadata.getUuid()); | ||
| } | ||
| pendingChunkedMessageCount++; | ||
| if (maxPendingChunkedMessage > 0 && pendingChunkedMessageCount > maxPendingChunkedMessage) { | ||
| removeOldestPendingChunkedMessage(); | ||
| pendingChunkedMessageUuidQueue.remove(msgMetadata.getUuid()); | ||
| pendingChunkedMessageUuidQueue.add(msgMetadata.getUuid()); | ||
| } else { | ||
| // Genuinely new uuid: count it and enqueue it exactly once. Eviction is only checked | ||
| // here because only a new uuid grows the number of in-flight chunked messages. | ||
| pendingChunkedMessageCount++; | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. [SUGGESTION] consider deriving the pending count from chunkedMessagesMap instead of a separate counter The bug here is that a hand-maintained If you keep the counter, that is fine too; it just means every future change to these paths has to keep it in step by hand. |
||
| if (maxPendingChunkedMessage > 0 && pendingChunkedMessageCount > maxPendingChunkedMessage) { | ||
| removeOldestPendingChunkedMessage(); | ||
| } | ||
| pendingChunkedMessageUuidQueue.add(msgMetadata.getUuid()); | ||
| } | ||
| int totalChunks = msgMetadata.getNumChunksFromMsg(); | ||
| ByteBuf chunkedMsgBuffer = PulsarByteBufAllocator.DEFAULT.buffer(msgMetadata.getTotalChunkMsgSize(), | ||
| msgMetadata.getTotalChunkMsgSize()); | ||
| chunkedMsgCtx = chunkedMessagesMap.computeIfAbsent(msgMetadata.getUuid(), | ||
| (key) -> ChunkedMessageCtx.get(totalChunks, chunkedMsgBuffer)); | ||
| pendingChunkedMessageUuidQueue.add(msgMetadata.getUuid()); | ||
| } | ||
|
|
||
| // discard message if chunk is out-of-order | ||
|
|
@@ -1692,6 +1721,14 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m | |
| ReferenceCountUtil.safeRelease(chunkedMsgCtx.chunkedMsgBuffer); | ||
| } | ||
| chunkedMsgCtx.recycle(); | ||
| // A non-null context here means an out-of-order chunk is discarding a previously | ||
| // tracked assembly. That assembly's uuid was counted AND added to | ||
| // pendingChunkedMessageUuidQueue when its first chunk arrived, so remove it from both | ||
| // to keep chunkedMessagesMap, pendingChunkedMessageCount and the uuid queue in sync. | ||
| // (When chunkedMsgCtx is null nothing was ever tracked for this uuid, so there is | ||
| // nothing to decrement or dequeue.) | ||
| pendingChunkedMessageCount--; | ||
| pendingChunkedMessageUuidQueue.remove(msgMetadata.getUuid()); | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. [NIT] queue.remove(uuid) is a linear scan on a locked array queue
|
||
| } | ||
| chunkedMessagesMap.remove(msgMetadata.getUuid()); | ||
| compressedPayload.release(); | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
[SUGGESTION] this test can assert before the duplicate chunks are processed, and the out-of-order path is not covered
producer.send()returning does not mean the consumer has processed the chunk, andchunkedMessagesMap.size() == 1is already true after the first of the ten duplicates, so the assertions below can run before the rest have arrived (and between a replacement'sremoveandcomputeIfAbsentthe map is briefly empty). It failed reliably for me with the fix reverted, so it is not vacuous in practice, but it is timing-dependent. A deterministic way is to send a complete second chunked message afterwards, receive it (messages are processed in order, so all duplicates are done), and then assert the three sizes.It would also be good to cover the out-of-order discard (
count--plusqueue.remove), and eviction order after a replacement.