From 75e33a5058e79a484c71633b1f0a2dc52278417b Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Thu, 31 Aug 2023 17:47:43 +0800 Subject: [PATCH 1/4] [fix][test]Flaky test testMaxPendingChunkMessages --- .../client/impl/MessageChunkingTest.java | 26 +++++++++++-------- 1 file changed, 15 insertions(+), 11 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..2957a2cd0a564 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 @@ -62,6 +62,7 @@ import org.apache.pulsar.common.protocol.Commands; import org.apache.pulsar.common.protocol.Commands.ChecksumType; import org.apache.pulsar.common.util.FutureUtil; +import org.awaitility.Awaitility; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.testng.Assert; @@ -339,21 +340,24 @@ public void testMaxPendingChunkMessages() throws Exception { .enableBatching(false) .create(); - sendSingleChunk(producer, "0", 0, 2); - sendSingleChunk(producer, "1", 0, 2); - sendSingleChunk(producer, "1", 1, 2); + Awaitility.await().untilAsserted(() -> { + sendSingleChunk(producer, "0", 0, 2); + sendSingleChunk(producer, "1", 0, 2); + sendSingleChunk(producer, "1", 1, 2); - // The chunked message of uuid 0 is discarded. - Message receivedMsg = consumer.receive(5, TimeUnit.SECONDS); - assertEquals(receivedMsg.getValue(), "chunk-1-0|chunk-1-1|"); + // The chunked message of uuid 0 is discarded. + Message receivedMsg = consumer.receive(5, TimeUnit.SECONDS); + assertEquals(receivedMsg.getValue(), "chunk-1-0|chunk-1-1|"); - consumer.acknowledge(receivedMsg); - consumer.redeliverUnacknowledgedMessages(); + consumer.acknowledge(receivedMsg); + consumer.redeliverUnacknowledgedMessages(); - sendSingleChunk(producer, "0", 1, 2); + sendSingleChunk(producer, "0", 1, 2); + + // Ensure that the chunked message of uuid 0 is discarded. + assertNull(consumer.receive(5, TimeUnit.SECONDS)); + }); - // Ensure that the chunked message of uuid 0 is discarded. - assertNull(consumer.receive(5, TimeUnit.SECONDS)); } @Test From 10056d4ae3af3339cf678a46db3674374743900e Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Thu, 31 Aug 2023 20:24:26 +0800 Subject: [PATCH 2/4] optimize --- .../org/apache/pulsar/client/impl/MessageChunkingTest.java | 6 +++++- .../java/org/apache/pulsar/client/impl/ConsumerImpl.java | 7 ++++--- 2 files changed, 9 insertions(+), 4 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 2957a2cd0a564..e9635314b3706 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 @@ -355,7 +355,11 @@ public void testMaxPendingChunkMessages() throws Exception { sendSingleChunk(producer, "0", 1, 2); // Ensure that the chunked message of uuid 0 is discarded. - assertNull(consumer.receive(5, TimeUnit.SECONDS)); + Message msg = consumer.receive(5, TimeUnit.SECONDS); + if (msg != null) { + consumer.acknowledge(msg); + } + assertNull(msg); }); } 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..45ec31c893f84 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 @@ -1486,9 +1486,10 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m return null; } // means we lost the first chunk: should never happen - log.info("[{}] [{}] Received unexpected chunk messageId {}, last-chunk-id = {}, chunkId = {}", topic, - subscription, msgId, - (chunkedMsgCtx != null ? chunkedMsgCtx.lastChunkedMessageId : null), msgMetadata.getChunkId()); + log.info("[{}] [{}] Received unexpected chunk messageId {}, last-chunk-id = {}, chunkId = {}, uuid = {}", + topic, subscription, msgId, + (chunkedMsgCtx != null ? chunkedMsgCtx.lastChunkedMessageId : null), msgMetadata.getChunkId(), + msgMetadata.getUuid()); if (chunkedMsgCtx != null) { if (chunkedMsgCtx.chunkedMsgBuffer != null) { ReferenceCountUtil.safeRelease(chunkedMsgCtx.chunkedMsgBuffer); From a5e90553bb0f76db0494ebdaaf14d761f23e4fd2 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Thu, 7 Sep 2023 09:12:11 +0800 Subject: [PATCH 3/4] add some note and optimize the test --- .../client/impl/MessageChunkingTest.java | 51 +++++++++++-------- 1 file changed, 31 insertions(+), 20 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 e9635314b3706..145e915d8ac19 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 @@ -320,15 +320,29 @@ private void sendSingleChunk(Producer producer, String uuid, int chunkId msg.send(); } + /** + * 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 + * consumer. + * For example: + * ChunkMessage1 chunk-1: uuid = 0, chunkId = 0, totalChunk = 2; + * ChunkMessage2 chunk-1: uuid = 1, chunkId = 0, totalChunk = 2; + * ChunkMessage2 chunk-2: uuid = 1, chunkId = 1, totalChunk = 2; + * ChunkMessage1 chunk-2: uuid = 0, chunkId = 1, totalChunk = 2; + * The chunk-1 in the ChunkMessage1 and ChunkMessage all is incomplete. + * chunk-1 in the ChunkMessage1 will be discarded and acked when receive the chunk-1 in the ChunkMessage2. + * If ack ChunkMessage2 and redeliver unacknowledged messages, the consumer can not receive any message again. + * @throws Exception + */ @Test public void testMaxPendingChunkMessages() throws Exception { log.info("-- Starting {} test --", methodName); final String topicName = "persistent://my-property/my-ns/maxPending"; - + final String subName = "my-subscriber-name"; @Cleanup Consumer consumer = pulsarClient.newConsumer(Schema.STRING) .topic(topicName) - .subscriptionName("my-subscriber-name") + .subscriptionName(subName) .maxPendingChunkedMessage(1) .autoAckOldestChunkedMessageOnQueueFull(true) .subscribe(); @@ -340,28 +354,25 @@ public void testMaxPendingChunkMessages() throws Exception { .enableBatching(false) .create(); - Awaitility.await().untilAsserted(() -> { - sendSingleChunk(producer, "0", 0, 2); - sendSingleChunk(producer, "1", 0, 2); - sendSingleChunk(producer, "1", 1, 2); - - // The chunked message of uuid 0 is discarded. - Message receivedMsg = consumer.receive(5, TimeUnit.SECONDS); - assertEquals(receivedMsg.getValue(), "chunk-1-0|chunk-1-1|"); + sendSingleChunk(producer, "0", 0, 2); + sendSingleChunk(producer, "1", 0, 2); + sendSingleChunk(producer, "1", 1, 2); - consumer.acknowledge(receivedMsg); - consumer.redeliverUnacknowledgedMessages(); + // The chunked message of uuid 0 is discarded. + Message receivedMsg = consumer.receive(5, TimeUnit.SECONDS); + assertEquals(receivedMsg.getValue(), "chunk-1-0|chunk-1-1|"); - sendSingleChunk(producer, "0", 1, 2); + consumer.acknowledge(receivedMsg); + assertEquals(admin.topics().getStats(topicName).getSubscriptions().get(subName) + .getNonContiguousDeletedMessagesRanges(), 0); + Awaitility.await().untilAsserted(() -> assertEquals(admin.topics().getStats(topicName) + .getSubscriptions().get(subName).getNonContiguousDeletedMessagesRanges(), 0)); + consumer.redeliverUnacknowledgedMessages(); - // Ensure that the chunked message of uuid 0 is discarded. - Message msg = consumer.receive(5, TimeUnit.SECONDS); - if (msg != null) { - consumer.acknowledge(msg); - } - assertNull(msg); - }); + sendSingleChunk(producer, "0", 1, 2); + // Ensure that the chunked message of uuid 0 is discarded. + assertNull(consumer.receive(5, TimeUnit.SECONDS)); } @Test From 41a56cfd6f12b4601a382f9772b6d96cf8a9c778 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Thu, 7 Sep 2023 09:31:55 +0800 Subject: [PATCH 4/4] fix --- .../org/apache/pulsar/client/impl/MessageChunkingTest.java | 4 +--- 1 file changed, 1 insertion(+), 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 145e915d8ac19..49b96cf79198b 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 @@ -329,7 +329,7 @@ private void sendSingleChunk(Producer producer, String uuid, int chunkId * ChunkMessage2 chunk-1: uuid = 1, chunkId = 0, totalChunk = 2; * ChunkMessage2 chunk-2: uuid = 1, chunkId = 1, totalChunk = 2; * ChunkMessage1 chunk-2: uuid = 0, chunkId = 1, totalChunk = 2; - * The chunk-1 in the ChunkMessage1 and ChunkMessage all is incomplete. + * The chunk-1 in the ChunkMessage1 and ChunkMessage2 all is incomplete. * chunk-1 in the ChunkMessage1 will be discarded and acked when receive the chunk-1 in the ChunkMessage2. * If ack ChunkMessage2 and redeliver unacknowledged messages, the consumer can not receive any message again. * @throws Exception @@ -363,8 +363,6 @@ public void testMaxPendingChunkMessages() throws Exception { assertEquals(receivedMsg.getValue(), "chunk-1-0|chunk-1-1|"); consumer.acknowledge(receivedMsg); - assertEquals(admin.topics().getStats(topicName).getSubscriptions().get(subName) - .getNonContiguousDeletedMessagesRanges(), 0); Awaitility.await().untilAsserted(() -> assertEquals(admin.topics().getStats(topicName) .getSubscriptions().get(subName).getNonContiguousDeletedMessagesRanges(), 0)); consumer.redeliverUnacknowledgedMessages();