From 1ba8188b5d3f01cb8dd3734f3259f4f6e98b1dca Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Thu, 31 Aug 2023 14:06:04 +0800 Subject: [PATCH 1/5] [fix][client] Avoid ack hole for chunk message ## Motivation Handle ack hole case: For example: Chunk-1 sequence ID: 0, chunk ID: 0, msgID: 1:1 Chunk-2 sequence ID: 0, chunk ID: 1, msgID: 1:2 Chunk-3 sequence ID: 0, chunk ID: 0, msgID: 1:3 Chunk-4 sequence ID: 0, chunk ID: 1, msgID: 1:4 Chunk-5 sequence ID: 0, chunk ID: 2, msgID: 1:5 Consumer ack chunk message via ChunkMessageIdImpl that is consist of all the chunks in this chunk message(Chunk-3, Chunk-4, Chunk-5). The Chunk-1 and Chunk-2 are not included in the ChunkMessageIdImpl, so we should process here. ## Modification Ack chunk-1 and chunk-2. --- .../pulsar/client/impl/ConsumerImpl.java | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) 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 ef063c6b15970..68f741ea07935 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 @@ -1437,6 +1437,24 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m if (msgMetadata.getChunkId() == 0) { if (chunkedMsgCtx != null) { + // Handle ack hole case: + // For example: + // Chunk-1 sequence ID: 0, chunk ID: 0, msgID: 1:1 + // Chunk-2 sequence ID: 0, chunk ID: 1, msgID: 1:2 + // Chunk-3 sequence ID: 0, chunk ID: 0, msgID: 1:3 + // Chunk-4 sequence ID: 0, chunk ID: 1, msgID: 1:4 + // Chunk-5 sequence ID: 0, chunk ID: 2, msgID: 1:5 + // Consumer ack chunk message via ChunkMessageIdImpl that is consist of all the chunks in this chunk + // message(Chunk-3, Chunk-4, Chunk-5). The Chunk-1 and Chunk-2 are not included in the + // ChunkMessageIdImpl, so we should process here. + boolean repeatedlyReceived = Arrays.stream(chunkedMsgCtx.chunkedMessageIds) + .anyMatch(messageId1 -> messageId1 != null && messageId1.ledgerId == messageId.getLedgerId() + && messageId1.entryId == messageId.getEntryId()); + if (!repeatedlyReceived) { + Arrays.stream(chunkedMsgCtx.chunkedMessageIds).forEach(messageId1 -> { + doAcknowledge(messageId1, AckType.Individual, Collections.emptyMap(), null); + }); + } // 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 From dbd763680aaed26d8489faf505dcba089df62c8d Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Thu, 31 Aug 2023 15:03:24 +0800 Subject: [PATCH 2/5] add test --- .../client/impl/MessageChunkingTest.java | 32 +++++++++++++++++++ .../pulsar/client/impl/ConsumerImpl.java | 8 +++-- 2 files changed, 37 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 dffa003524864..02130339e3b7f 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 @@ -356,6 +356,38 @@ public void testMaxPendingChunkMessages() throws Exception { assertNull(consumer.receive(5, TimeUnit.SECONDS)); } + @Test + public void testResendChunkMessagesWithoutAckHole() throws Exception { + log.info("-- Starting {} test --", methodName); + final String topicName = "persistent://my-property/my-ns/testResendChunkMessagesWithoutAckHole"; + final String subName = "my-subscriber-name"; + @Cleanup + Consumer consumer = pulsarClient.newConsumer(Schema.STRING) + .topic(topicName) + .subscriptionName(subName) + .maxPendingChunkedMessage(10) + .autoAckOldestChunkedMessageOnQueueFull(true) + .subscribe(); + @Cleanup + Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topicName) + .chunkMaxMessageSize(100) + .enableChunking(true) + .enableBatching(false) + .create(); + + sendSingleChunk(producer, "0", 0, 2); + + sendSingleChunk(producer, "0", 0, 2); // Resending the first chunk + sendSingleChunk(producer, "0", 1, 2); + + Message receivedMsg = consumer.receive(5, TimeUnit.SECONDS); + assertEquals(receivedMsg.getValue(), "chunk-0-0|chunk-0-1|"); + consumer.acknowledge(receivedMsg); + assertEquals(admin.topics().getStats(topicName).getSubscriptions().get(subName) + .getNonContiguousDeletedMessagesRanges(), 0); + } + @Test public void testResendChunkMessages() throws Exception { log.info("-- Starting {} test --", methodName); 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 68f741ea07935..75aed15591169 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 @@ -1444,15 +1444,17 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m // Chunk-3 sequence ID: 0, chunk ID: 0, msgID: 1:3 // Chunk-4 sequence ID: 0, chunk ID: 1, msgID: 1:4 // Chunk-5 sequence ID: 0, chunk ID: 2, msgID: 1:5 - // Consumer ack chunk message via ChunkMessageIdImpl that is consist of all the chunks in this chunk + // Consumer ack chunk message via ChunkMessageIdImpl that consists of all the chunks in this chunk // message(Chunk-3, Chunk-4, Chunk-5). The Chunk-1 and Chunk-2 are not included in the - // ChunkMessageIdImpl, so we should process here. + // ChunkMessageIdImpl, so we should process it here. boolean repeatedlyReceived = Arrays.stream(chunkedMsgCtx.chunkedMessageIds) .anyMatch(messageId1 -> messageId1 != null && messageId1.ledgerId == messageId.getLedgerId() && messageId1.entryId == messageId.getEntryId()); if (!repeatedlyReceived) { Arrays.stream(chunkedMsgCtx.chunkedMessageIds).forEach(messageId1 -> { - doAcknowledge(messageId1, AckType.Individual, Collections.emptyMap(), null); + if (messageId1 != null) { + doAcknowledge(messageId1, AckType.Individual, Collections.emptyMap(), null); + } }); } // The first chunk of a new chunked-message received before receiving other chunks of previous From f6595f87d77c4bdf064b7c48a3d0374e69784b9f Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Thu, 31 Aug 2023 23:27:03 +0800 Subject: [PATCH 3/5] optimize --- .../client/impl/MessageChunkingTest.java | 1 + .../pulsar/client/impl/ConsumerImpl.java | 34 ++++++++++++++----- 2 files changed, 26 insertions(+), 9 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 02130339e3b7f..136585c0df238 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 @@ -427,6 +427,7 @@ public void testResendChunkMessages() throws Exception { receivedMsg = consumer.receive(5, TimeUnit.SECONDS); assertEquals(receivedMsg.getValue(), "chunk-1-0|chunk-1-1|chunk-1-2|"); consumer.acknowledge(receivedMsg); + Assert.assertEquals(((ConsumerImpl) consumer).getAvailablePermits(), 10); } /** 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 75aed15591169..92c6829cffd1f 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 @@ -1437,20 +1437,33 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m if (msgMetadata.getChunkId() == 0) { if (chunkedMsgCtx != null) { - // Handle ack hole case: + // Handle ack hole case when receive duplicated chunks. + // There are two situation that receives chunks with the same sequence ID and chunk ID. + // Situation 1 - Message redeliver: + // For example: + // Chunk-1 sequence ID: 0, chunk ID: 0, msgID: 1:1 + // Chunk-2 sequence ID: 0, chunk ID: 1, msgID: 1:2 + // Chunk-3 sequence ID: 0, chunk ID: 0, msgID: 1:1 + // Chunk-4 sequence ID: 0, chunk ID: 1, msgID: 1:2 + // Chunk-5 sequence ID: 0, chunk ID: 2, msgID: 1:3 + // In this case, chunk-3 and chunk-4 have the same msgID with chunk-1 and chunk-2. + // This may be caused by message redeliver, we can't ack any chunk in this case here. + // Situation 2 - Message duplication: // For example: // Chunk-1 sequence ID: 0, chunk ID: 0, msgID: 1:1 // Chunk-2 sequence ID: 0, chunk ID: 1, msgID: 1:2 // Chunk-3 sequence ID: 0, chunk ID: 0, msgID: 1:3 // Chunk-4 sequence ID: 0, chunk ID: 1, msgID: 1:4 // Chunk-5 sequence ID: 0, chunk ID: 2, msgID: 1:5 + // In this case, all the chunks have the different msgID. Chunk-1, Chunk-2, Chunk-3, Chunk-4 are + // duplicated persisting in the topic. // Consumer ack chunk message via ChunkMessageIdImpl that consists of all the chunks in this chunk - // message(Chunk-3, Chunk-4, Chunk-5). The Chunk-1 and Chunk-2 are not included in the - // ChunkMessageIdImpl, so we should process it here. - boolean repeatedlyReceived = Arrays.stream(chunkedMsgCtx.chunkedMessageIds) - .anyMatch(messageId1 -> messageId1 != null && messageId1.ledgerId == messageId.getLedgerId() + // message(Chunk-3, Chunk-4, Chunk-5). The Chunk-1 and Chunk-2 would not be included in the + // ChunkMessageIdImpl, so we should ack them here to avoid ack hole. + boolean messageDuplication = Arrays.stream(chunkedMsgCtx.chunkedMessageIds) + .noneMatch(messageId1 -> messageId1 != null && messageId1.ledgerId == messageId.getLedgerId() && messageId1.entryId == messageId.getEntryId()); - if (!repeatedlyReceived) { + if (messageDuplication) { Arrays.stream(chunkedMsgCtx.chunkedMessageIds).forEach(messageId1 -> { if (messageId1 != null) { doAcknowledge(messageId1, AckType.Individual, Collections.emptyMap(), null); @@ -1465,6 +1478,7 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m } chunkedMsgCtx.recycle(); chunkedMessagesMap.remove(msgMetadata.getUuid()); + increaseAvailablePermits(cnx); } pendingChunkedMessageCount++; if (maxPendingChunkedMessage > 0 && pendingChunkedMessageCount > maxPendingChunkedMessage) { @@ -1497,10 +1511,12 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m msgMetadata.getChunkId(), msgMetadata.getSequenceId()); compressedPayload.release(); increaseAvailablePermits(cnx); - boolean repeatedlyReceived = Arrays.stream(chunkedMsgCtx.chunkedMessageIds) - .anyMatch(messageId1 -> messageId1 != null && messageId1.ledgerId == messageId.getLedgerId() + // Just like the above logic of receiving the first chunk again. We only ack this chunk in the message + // duplication case. + boolean messageDuplication = Arrays.stream(chunkedMsgCtx.chunkedMessageIds) + .noneMatch(messageId1 -> messageId1 != null && messageId1.ledgerId == messageId.getLedgerId() && messageId1.entryId == messageId.getEntryId()); - if (!repeatedlyReceived) { + if (messageDuplication) { doAcknowledge(msgId, AckType.Individual, Collections.emptyMap(), null); } return null; From 48a7ee0c75535c96bef2cf8e04ffff3a4e6d19da Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Fri, 1 Sep 2023 19:54:07 +0800 Subject: [PATCH 4/5] optimize --- .../pulsar/client/impl/ConsumerImpl.java | 18 ++++++++---------- 1 file changed, 8 insertions(+), 10 deletions(-) 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 92c6829cffd1f..761c72b9ff87d 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 @@ -1448,22 +1448,20 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m // Chunk-5 sequence ID: 0, chunk ID: 2, msgID: 1:3 // In this case, chunk-3 and chunk-4 have the same msgID with chunk-1 and chunk-2. // This may be caused by message redeliver, we can't ack any chunk in this case here. - // Situation 2 - Message duplication: + // Situation 2 - Corrupted chunk message // For example: // Chunk-1 sequence ID: 0, chunk ID: 0, msgID: 1:1 // Chunk-2 sequence ID: 0, chunk ID: 1, msgID: 1:2 // Chunk-3 sequence ID: 0, chunk ID: 0, msgID: 1:3 // Chunk-4 sequence ID: 0, chunk ID: 1, msgID: 1:4 // Chunk-5 sequence ID: 0, chunk ID: 2, msgID: 1:5 - // In this case, all the chunks have the different msgID. Chunk-1, Chunk-2, Chunk-3, Chunk-4 are - // duplicated persisting in the topic. - // Consumer ack chunk message via ChunkMessageIdImpl that consists of all the chunks in this chunk - // message(Chunk-3, Chunk-4, Chunk-5). The Chunk-1 and Chunk-2 would not be included in the - // ChunkMessageIdImpl, so we should ack them here to avoid ack hole. - boolean messageDuplication = Arrays.stream(chunkedMsgCtx.chunkedMessageIds) + // In this case, all the chunks with different msgIDs and are persistent in the topic. + // But Chunk-1 and Chunk-2 belong to a corrupted chunk message that must be skipped since + // they will not be delivered to end users. So we should ack them here to avoid ack hole. + boolean isCorruptedChunkMessageDetected = Arrays.stream(chunkedMsgCtx.chunkedMessageIds) .noneMatch(messageId1 -> messageId1 != null && messageId1.ledgerId == messageId.getLedgerId() && messageId1.entryId == messageId.getEntryId()); - if (messageDuplication) { + if (isCorruptedChunkMessageDetected) { Arrays.stream(chunkedMsgCtx.chunkedMessageIds).forEach(messageId1 -> { if (messageId1 != null) { doAcknowledge(messageId1, AckType.Individual, Collections.emptyMap(), null); @@ -1513,10 +1511,10 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m increaseAvailablePermits(cnx); // Just like the above logic of receiving the first chunk again. We only ack this chunk in the message // duplication case. - boolean messageDuplication = Arrays.stream(chunkedMsgCtx.chunkedMessageIds) + boolean isDuplicatedChunk = Arrays.stream(chunkedMsgCtx.chunkedMessageIds) .noneMatch(messageId1 -> messageId1 != null && messageId1.ledgerId == messageId.getLedgerId() && messageId1.entryId == messageId.getEntryId()); - if (messageDuplication) { + if (isDuplicatedChunk) { doAcknowledge(msgId, AckType.Individual, Collections.emptyMap(), null); } return null; From b45ca0a0ac5575264c467a6a3a08f94b04f8f814 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Fri, 1 Sep 2023 20:47:35 +0800 Subject: [PATCH 5/5] optimize --- .../apache/pulsar/client/impl/MessageChunkingTest.java | 2 +- .../java/org/apache/pulsar/client/impl/ConsumerImpl.java | 8 +++----- 2 files changed, 4 insertions(+), 6 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 136585c0df238..f266afd8a2ee1 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 @@ -427,7 +427,7 @@ public void testResendChunkMessages() throws Exception { receivedMsg = consumer.receive(5, TimeUnit.SECONDS); assertEquals(receivedMsg.getValue(), "chunk-1-0|chunk-1-1|chunk-1-2|"); consumer.acknowledge(receivedMsg); - Assert.assertEquals(((ConsumerImpl) consumer).getAvailablePermits(), 10); + Assert.assertEquals(((ConsumerImpl) consumer).getAvailablePermits(), 8); } /** 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 761c72b9ff87d..34286f51bfaf8 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 @@ -1421,7 +1421,9 @@ void messageReceived(CommandMessage cmdMessage, ByteBuf headersAndPayload, Clien private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata msgMetadata, MessageIdImpl msgId, MessageIdData messageId, ClientCnx cnx) { - + if (msgMetadata.getChunkId() != (msgMetadata.getNumChunksFromMsg() - 1)) { + increaseAvailablePermits(cnx); + } // Lazy task scheduling to expire incomplete chunk message if (expireTimeOfIncompleteChunkedMessageMillis > 0 && expireChunkMessageTaskScheduled.compareAndSet(false, true)) { @@ -1476,7 +1478,6 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m } chunkedMsgCtx.recycle(); chunkedMessagesMap.remove(msgMetadata.getUuid()); - increaseAvailablePermits(cnx); } pendingChunkedMessageCount++; if (maxPendingChunkedMessage > 0 && pendingChunkedMessageCount > maxPendingChunkedMessage) { @@ -1508,7 +1509,6 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m msgMetadata.getProducerName(), msgId, chunkedMsgCtx.lastChunkedMessageId, msgMetadata.getChunkId(), msgMetadata.getSequenceId()); compressedPayload.release(); - increaseAvailablePermits(cnx); // Just like the above logic of receiving the first chunk again. We only ack this chunk in the message // duplication case. boolean isDuplicatedChunk = Arrays.stream(chunkedMsgCtx.chunkedMessageIds) @@ -1531,7 +1531,6 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m } chunkedMessagesMap.remove(msgMetadata.getUuid()); compressedPayload.release(); - increaseAvailablePermits(cnx); if (expireTimeOfIncompleteChunkedMessageMillis > 0 && System.currentTimeMillis() > (msgMetadata.getPublishTime() + expireTimeOfIncompleteChunkedMessageMillis)) { @@ -1550,7 +1549,6 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m // if final chunk is not received yet then release payload and return if (msgMetadata.getChunkId() != (msgMetadata.getNumChunksFromMsg() - 1)) { compressedPayload.release(); - increaseAvailablePermits(cnx); return null; }