From 52ae60f919bcf6b668b50d47878bf4d18edb2c80 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Mon, 7 Aug 2023 23:33:32 +0800 Subject: [PATCH 01/24] [fix][broker]Fix chunked messages will be filtered by duplicating ### Motivation Chunked messages use the same metadata, so all the chunked messages in a single message use the same sequence Id. And it will be recorded as duplicated messages. ``` private long updateMessageMetadataSequenceId(final MessageMetadata msgMetadata) { final long sequenceId; if (!msgMetadata.hasSequenceId()) { sequenceId = msgIdGenerator++; msgMetadata.setSequenceId(sequenceId); } else { sequenceId = msgMetadata.getSequenceId(); } return sequenceId; } ``` ### Modification Use different sequence id for chunk message. --- .../impl/MessageChunkingSharedTest.java | 27 +++++++++++++++++++ .../pulsar/client/impl/ProducerImpl.java | 11 ++++---- 2 files changed, 33 insertions(+), 5 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java index 163c42d835b35..8e90b48f55178 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java @@ -42,6 +42,7 @@ import org.apache.pulsar.client.api.MessageListener; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.ProducerConsumerBase; +import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionType; @@ -211,6 +212,32 @@ private static String createChunkedMessage(int numChunks) { return Schema.STRING.decode(payload); } + @Test + public void testDuplicateForChunkMessage() throws Exception { + String topicName = "persistent://my-property/my-ns/testDuplicateForChunkMessage"; + String producerName = "test-producer"; + pulsarClient = PulsarClient.builder().serviceUrl(pulsar.getBrokerServiceUrl()).build(); + // consumer + Consumer consumer = pulsarClient + .newConsumer() + .subscriptionName("test-sub") + .topic(topicName) + .subscribe(); + // producer + Producer partProducer = pulsarClient + .newProducer(Schema.AVRO(String.class)) + .producerName(producerName) + .topic(topicName) + .enableChunking(true) + .enableBatching(false) + .create(); + int messageSize = 6000; // payload size in KB + String message = "a".repeat(messageSize * 1000); + partProducer.newMessage().value(message).send(); + Message msg = consumer.receive(5, TimeUnit.SECONDS); + assertNotNull(msg); + } + private static void sendNonChunk(final PersistentTopic persistentTopic, final String producerName, final long sequenceId) { diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java index 267b06649d719..8bc5d3613a411 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java @@ -51,6 +51,7 @@ import java.util.Map; import java.util.Optional; import java.util.Queue; +import java.util.UUID; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.Semaphore; @@ -531,6 +532,7 @@ public void sendAsync(Message message, SendCallback callback) { ? msg.getMessageBuilder().getOrderingKey() : null; // msg.messageId will be reset if previous message chunk is sent successfully. final MessageId messageId = msg.getMessageId(); + final UUID chunkUUid = UUID.randomUUID(); for (int chunkId = 0; chunkId < totalChunks; chunkId++) { // Need to reset the schemaVersion, because the schemaVersion is based on a ByteBuf object in // `MessageMetadata`, if we want to re-serialize the `SEND` command using a same `MessageMetadata`, @@ -553,9 +555,8 @@ public void sendAsync(Message message, SendCallback callback) { synchronized (this) { // Update the message metadata before computing the payload chunk size // to avoid a large message cannot be split into chunks. - final long sequenceId = updateMessageMetadataSequenceId(msgMetadata); - String uuid = totalChunks > 1 ? String.format("%s-%d", producerName, sequenceId) : null; - + final long sequenceId = updateMessageMetadataSequenceId(msgMetadata, chunkId > 0); + String uuid = totalChunks > 1 ? String.format("%s-%s", producerName, chunkUUid) : null; serializeAndSendMessage(msg, payload, sequenceId, uuid, chunkId, totalChunks, readStartIndex, payloadChunkSize, compressedPayload, compressed, compressedPayload.readableBytes(), callback, chunkedMessageCtx, messageId); @@ -594,9 +595,9 @@ private void updateMessageMetadata(final MessageMetadata msgMetadata, final int } } - private long updateMessageMetadataSequenceId(final MessageMetadata msgMetadata) { + private long updateMessageMetadataSequenceId(final MessageMetadata msgMetadata, boolean isChunk) { final long sequenceId; - if (!msgMetadata.hasSequenceId()) { + if (!msgMetadata.hasSequenceId() || isChunk) { sequenceId = msgIdGenerator++; msgMetadata.setSequenceId(sequenceId); } else { From 30c9430a86c0602e6baf113553761331fdae6e17 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Thu, 10 Aug 2023 10:00:51 +0800 Subject: [PATCH 02/24] optimize --- .../impl/MessageChunkingSharedTest.java | 24 +++++++++++++++---- .../pulsar/client/impl/ProducerImpl.java | 9 +++++-- 2 files changed, 27 insertions(+), 6 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java index 8e90b48f55178..05929f0733ad4 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java @@ -23,6 +23,7 @@ import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertTrue; import io.netty.buffer.ByteBuf; +import java.lang.reflect.Field; import java.time.Duration; import java.util.ArrayList; import java.util.Arrays; @@ -214,18 +215,20 @@ private static String createChunkedMessage(int numChunks) { @Test public void testDuplicateForChunkMessage() throws Exception { + this.conf.setBrokerDeduplicationEnabled(true); + restartBroker(); String topicName = "persistent://my-property/my-ns/testDuplicateForChunkMessage"; String producerName = "test-producer"; pulsarClient = PulsarClient.builder().serviceUrl(pulsar.getBrokerServiceUrl()).build(); // consumer - Consumer consumer = pulsarClient - .newConsumer() + Consumer consumer = pulsarClient + .newConsumer(Schema.STRING) .subscriptionName("test-sub") .topic(topicName) .subscribe(); // producer Producer partProducer = pulsarClient - .newProducer(Schema.AVRO(String.class)) + .newProducer(Schema.STRING) .producerName(producerName) .topic(topicName) .enableChunking(true) @@ -234,8 +237,21 @@ public void testDuplicateForChunkMessage() throws Exception { int messageSize = 6000; // payload size in KB String message = "a".repeat(messageSize * 1000); partProducer.newMessage().value(message).send(); - Message msg = consumer.receive(5, TimeUnit.SECONDS); + Message msg = consumer.receive(5, TimeUnit.SECONDS); assertNotNull(msg); + assertEquals(msg.getValue(), message); + + Field msgIdGenerator = ProducerImpl.class.getDeclaredField("msgIdGenerator"); + msgIdGenerator.setAccessible(true); + assertEquals(msg.getSequenceId() + 2, msgIdGenerator.get(partProducer)); + + long sequenceID = (long) msgIdGenerator.get(partProducer); + String message2 = "b".repeat(messageSize * 1000); + partProducer.newMessage().value(message2).sequenceId(sequenceID).send(); + Message msg2 = consumer.receive(5, TimeUnit.SECONDS); + assertNotNull(msg2); + assertEquals(msg2.getValue(), message2); + assertEquals(msg2.getSequenceId(), sequenceID); } private static void sendNonChunk(final PersistentTopic persistentTopic, diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java index 8bc5d3613a411..0167c569bd34e 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java @@ -556,7 +556,7 @@ public void sendAsync(Message message, SendCallback callback) { // Update the message metadata before computing the payload chunk size // to avoid a large message cannot be split into chunks. final long sequenceId = updateMessageMetadataSequenceId(msgMetadata, chunkId > 0); - String uuid = totalChunks > 1 ? String.format("%s-%s", producerName, chunkUUid) : null; + String uuid = totalChunks > 1 ? String.format("%s-%s", producerName, sequenceId - chunkId) : null; serializeAndSendMessage(msg, payload, sequenceId, uuid, chunkId, totalChunks, readStartIndex, payloadChunkSize, compressedPayload, compressed, compressedPayload.readableBytes(), callback, chunkedMessageCtx, messageId); @@ -597,11 +597,16 @@ private void updateMessageMetadata(final MessageMetadata msgMetadata, final int private long updateMessageMetadataSequenceId(final MessageMetadata msgMetadata, boolean isChunk) { final long sequenceId; - if (!msgMetadata.hasSequenceId() || isChunk) { + // We always need to increment the value of `msgIdGenerator`, + // regardless of whether the user has set a sequence ID when sending a message. + if (!msgMetadata.hasSequenceId()) { sequenceId = msgIdGenerator++; msgMetadata.setSequenceId(sequenceId); + } else if (isChunk) { + sequenceId = msgIdGenerator++; } else { sequenceId = msgMetadata.getSequenceId(); + msgIdGenerator++; } return sequenceId; } From 02d70712060e4b5fd69c351db703dd3e3931a158 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Thu, 10 Aug 2023 20:18:56 +0800 Subject: [PATCH 03/24] optimize --- .../apache/pulsar/client/impl/MessageChunkingSharedTest.java | 4 +++- .../main/java/org/apache/pulsar/client/impl/ProducerImpl.java | 3 +-- 2 files changed, 4 insertions(+), 3 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java index 05929f0733ad4..c4b3dd0121eb5 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java @@ -239,17 +239,19 @@ public void testDuplicateForChunkMessage() throws Exception { partProducer.newMessage().value(message).send(); Message msg = consumer.receive(5, TimeUnit.SECONDS); assertNotNull(msg); + assertTrue(msg.getMessageId() instanceof ChunkMessageIdImpl); assertEquals(msg.getValue(), message); Field msgIdGenerator = ProducerImpl.class.getDeclaredField("msgIdGenerator"); msgIdGenerator.setAccessible(true); assertEquals(msg.getSequenceId() + 2, msgIdGenerator.get(partProducer)); - long sequenceID = (long) msgIdGenerator.get(partProducer); + long sequenceID = (long) msgIdGenerator.get(partProducer) + 1024L; String message2 = "b".repeat(messageSize * 1000); partProducer.newMessage().value(message2).sequenceId(sequenceID).send(); Message msg2 = consumer.receive(5, TimeUnit.SECONDS); assertNotNull(msg2); + assertTrue(msg2.getMessageId() instanceof ChunkMessageIdImpl); assertEquals(msg2.getValue(), message2); assertEquals(msg2.getSequenceId(), sequenceID); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java index 0167c569bd34e..bec95f0636c6b 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java @@ -532,7 +532,6 @@ public void sendAsync(Message message, SendCallback callback) { ? msg.getMessageBuilder().getOrderingKey() : null; // msg.messageId will be reset if previous message chunk is sent successfully. final MessageId messageId = msg.getMessageId(); - final UUID chunkUUid = UUID.randomUUID(); for (int chunkId = 0; chunkId < totalChunks; chunkId++) { // Need to reset the schemaVersion, because the schemaVersion is based on a ByteBuf object in // `MessageMetadata`, if we want to re-serialize the `SEND` command using a same `MessageMetadata`, @@ -606,7 +605,7 @@ private long updateMessageMetadataSequenceId(final MessageMetadata msgMetadata, sequenceId = msgIdGenerator++; } else { sequenceId = msgMetadata.getSequenceId(); - msgIdGenerator++; + msgIdGenerator = msgMetadata.getSequenceId() + 1; } return sequenceId; } From d99ed6d908993c3183fb6217ff68e77b205b6829 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Fri, 11 Aug 2023 11:17:08 +0800 Subject: [PATCH 04/24] handle in broker side --- .../persistent/MessageDeduplication.java | 19 +++++++++++++++++++ .../impl/MessageChunkingSharedTest.java | 2 +- .../pulsar/client/impl/ProducerImpl.java | 13 ++++--------- 3 files changed, 24 insertions(+), 10 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java index ed4e70bfd2953..99934c67aa215 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java @@ -55,6 +55,8 @@ public class MessageDeduplication { private final ManagedLedger managedLedger; private ManagedCursor managedCursor; + private final ConcurrentOpenHashMap chunkMessageOngoing; + enum Status { // Deduplication is initialized @@ -139,6 +141,7 @@ public MessageDeduplication(PulsarService pulsar, PersistentTopic topic, Managed this.maxNumberOfProducers = pulsar.getConfiguration().getBrokerDeduplicationMaxNumberOfProducers(); this.snapshotCounter = 0; this.replicatorPrefix = pulsar.getConfiguration().getReplicatorPrefix(); + this.chunkMessageOngoing = ConcurrentOpenHashMap.newBuilder().build(); } private CompletableFuture recoverSequenceIdsMap() { @@ -323,6 +326,19 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade String producerName = publishContext.getProducerName(); long sequenceId = publishContext.getSequenceId(); + // The process of the Producer sending chunk messages is continuous, and all chunks of a message use the same + // message metadata and sequence ID. Therefore, it is only necessary to check if the sequence ID of the first + // chunk is duplicated. + // When we receive the initial message of a non-duplicated chunk message, we place it in the + // chunkMessageOngoing. Upon completion of sending this chunk message, if we receive other messages + // sent by this Producer, we will remove it from the chunkMessageOngoing. + if (chunkMessageOngoing.containsKey(producerName)) { + if (publishContext.isChunked() && chunkMessageOngoing.get(producerName).equals(sequenceId)) { + return MessageDupStatus.NotDup; + } else { + chunkMessageOngoing.remove(producerName); + } + } long highestSequenceId = Math.max(publishContext.getHighestSequenceId(), sequenceId); if (producerName.startsWith(replicatorPrefix)) { // Message is coming from replication, we need to use the original producer name and sequence id @@ -363,6 +379,9 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade } highestSequencedPushed.put(producerName, highestSequenceId); } + if (publishContext.isChunked()) { + chunkMessageOngoing.put(producerName, sequenceId); + } return MessageDupStatus.NotDup; } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java index c4b3dd0121eb5..cbc3ea68f0648 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java @@ -244,7 +244,7 @@ public void testDuplicateForChunkMessage() throws Exception { Field msgIdGenerator = ProducerImpl.class.getDeclaredField("msgIdGenerator"); msgIdGenerator.setAccessible(true); - assertEquals(msg.getSequenceId() + 2, msgIdGenerator.get(partProducer)); + assertEquals(msg.getSequenceId() + 1, msgIdGenerator.get(partProducer)); long sequenceID = (long) msgIdGenerator.get(partProducer) + 1024L; String message2 = "b".repeat(messageSize * 1000); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java index bec95f0636c6b..267b06649d719 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java @@ -51,7 +51,6 @@ import java.util.Map; import java.util.Optional; import java.util.Queue; -import java.util.UUID; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.Semaphore; @@ -554,8 +553,9 @@ public void sendAsync(Message message, SendCallback callback) { synchronized (this) { // Update the message metadata before computing the payload chunk size // to avoid a large message cannot be split into chunks. - final long sequenceId = updateMessageMetadataSequenceId(msgMetadata, chunkId > 0); - String uuid = totalChunks > 1 ? String.format("%s-%s", producerName, sequenceId - chunkId) : null; + final long sequenceId = updateMessageMetadataSequenceId(msgMetadata); + String uuid = totalChunks > 1 ? String.format("%s-%d", producerName, sequenceId) : null; + serializeAndSendMessage(msg, payload, sequenceId, uuid, chunkId, totalChunks, readStartIndex, payloadChunkSize, compressedPayload, compressed, compressedPayload.readableBytes(), callback, chunkedMessageCtx, messageId); @@ -594,18 +594,13 @@ private void updateMessageMetadata(final MessageMetadata msgMetadata, final int } } - private long updateMessageMetadataSequenceId(final MessageMetadata msgMetadata, boolean isChunk) { + private long updateMessageMetadataSequenceId(final MessageMetadata msgMetadata) { final long sequenceId; - // We always need to increment the value of `msgIdGenerator`, - // regardless of whether the user has set a sequence ID when sending a message. if (!msgMetadata.hasSequenceId()) { sequenceId = msgIdGenerator++; msgMetadata.setSequenceId(sequenceId); - } else if (isChunk) { - sequenceId = msgIdGenerator++; } else { sequenceId = msgMetadata.getSequenceId(); - msgIdGenerator = msgMetadata.getSequenceId() + 1; } return sequenceId; } From be88b7bac931d7e2098b317507c8774199a63faa Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Fri, 11 Aug 2023 18:36:26 +0800 Subject: [PATCH 05/24] broker reload and message resend --- .../persistent/MessageDeduplication.java | 28 ++++++++++++------- 1 file changed, 18 insertions(+), 10 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java index 99934c67aa215..be08c6a7d3e51 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java @@ -55,7 +55,7 @@ public class MessageDeduplication { private final ManagedLedger managedLedger; private ManagedCursor managedCursor; - private final ConcurrentOpenHashMap chunkMessageOngoing; + private final ConcurrentOpenHashMap chunkMessageOngoing; enum Status { @@ -141,7 +141,7 @@ public MessageDeduplication(PulsarService pulsar, PersistentTopic topic, Managed this.maxNumberOfProducers = pulsar.getConfiguration().getBrokerDeduplicationMaxNumberOfProducers(); this.snapshotCounter = 0; this.replicatorPrefix = pulsar.getConfiguration().getReplicatorPrefix(); - this.chunkMessageOngoing = ConcurrentOpenHashMap.newBuilder().build(); + this.chunkMessageOngoing = ConcurrentOpenHashMap.newBuilder().build(); } private CompletableFuture recoverSequenceIdsMap() { @@ -178,6 +178,12 @@ public void readEntriesComplete(List entries, Object ctx) { long sequenceId = Math.max(md.getHighestSequenceId(), md.getSequenceId()); highestSequencedPushed.put(producerName, sequenceId); highestSequencedPersisted.put(producerName, sequenceId); + // Only maintain the latest ongoing chunk for each producer in the chunkMessageOngoing. + if (md.hasChunkId()) { + chunkMessageOngoing.put(producerName, md.getChunkId()); + } else { + chunkMessageOngoing.remove(producerName); + } producerRemoved(producerName); entry.release(); @@ -332,8 +338,13 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade // When we receive the initial message of a non-duplicated chunk message, we place it in the // chunkMessageOngoing. Upon completion of sending this chunk message, if we receive other messages // sent by this Producer, we will remove it from the chunkMessageOngoing. + headersAndPayload.markReaderIndex(); + MessageMetadata msgMetadata = Commands.parseMessageMetadata(headersAndPayload); + headersAndPayload.resetReaderIndex(); if (chunkMessageOngoing.containsKey(producerName)) { - if (publishContext.isChunked() && chunkMessageOngoing.get(producerName).equals(sequenceId)) { + //Deduplication requires the dependency on the ordered nature of messages to work accurately. + if (publishContext.isChunked() && chunkMessageOngoing.get(producerName) < msgMetadata.getChunkId()) { + chunkMessageOngoing.put(producerName, msgMetadata.getChunkId()); return MessageDupStatus.NotDup; } else { chunkMessageOngoing.remove(producerName); @@ -343,15 +354,12 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade if (producerName.startsWith(replicatorPrefix)) { // Message is coming from replication, we need to use the original producer name and sequence id // for the purpose of deduplication and not rely on the "replicator" name. - int readerIndex = headersAndPayload.readerIndex(); - MessageMetadata md = Commands.parseMessageMetadata(headersAndPayload); - producerName = md.getProducerName(); - sequenceId = md.getSequenceId(); - highestSequenceId = Math.max(md.getHighestSequenceId(), sequenceId); + producerName = msgMetadata.getProducerName(); + sequenceId = msgMetadata.getSequenceId(); + highestSequenceId = Math.max(msgMetadata.getHighestSequenceId(), sequenceId); publishContext.setOriginalProducerName(producerName); publishContext.setOriginalSequenceId(sequenceId); publishContext.setOriginalHighestSequenceId(highestSequenceId); - headersAndPayload.readerIndex(readerIndex); } // Synchronize the get() and subsequent put() on the map. This would only be relevant if the producer @@ -380,7 +388,7 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade highestSequencedPushed.put(producerName, highestSequenceId); } if (publishContext.isChunked()) { - chunkMessageOngoing.put(producerName, sequenceId); + chunkMessageOngoing.put(producerName, msgMetadata.getChunkId()); } return MessageDupStatus.NotDup; } From 1c8eec371156eef710d603f61a5b6f8e2614fe3f Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Sat, 12 Aug 2023 16:01:23 +0800 Subject: [PATCH 06/24] Only check sequence ID --- .../persistent/MessageDeduplication.java | 37 ++++--------------- .../impl/MessageChunkingSharedTest.java | 18 ++++++--- 2 files changed, 19 insertions(+), 36 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java index be08c6a7d3e51..d51ab6b8eb07f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java @@ -55,8 +55,6 @@ public class MessageDeduplication { private final ManagedLedger managedLedger; private ManagedCursor managedCursor; - private final ConcurrentOpenHashMap chunkMessageOngoing; - enum Status { // Deduplication is initialized @@ -141,7 +139,6 @@ public MessageDeduplication(PulsarService pulsar, PersistentTopic topic, Managed this.maxNumberOfProducers = pulsar.getConfiguration().getBrokerDeduplicationMaxNumberOfProducers(); this.snapshotCounter = 0; this.replicatorPrefix = pulsar.getConfiguration().getReplicatorPrefix(); - this.chunkMessageOngoing = ConcurrentOpenHashMap.newBuilder().build(); } private CompletableFuture recoverSequenceIdsMap() { @@ -178,12 +175,6 @@ public void readEntriesComplete(List entries, Object ctx) { long sequenceId = Math.max(md.getHighestSequenceId(), md.getSequenceId()); highestSequencedPushed.put(producerName, sequenceId); highestSequencedPersisted.put(producerName, sequenceId); - // Only maintain the latest ongoing chunk for each producer in the chunkMessageOngoing. - if (md.hasChunkId()) { - chunkMessageOngoing.put(producerName, md.getChunkId()); - } else { - chunkMessageOngoing.remove(producerName); - } producerRemoved(producerName); entry.release(); @@ -332,24 +323,9 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade String producerName = publishContext.getProducerName(); long sequenceId = publishContext.getSequenceId(); - // The process of the Producer sending chunk messages is continuous, and all chunks of a message use the same - // message metadata and sequence ID. Therefore, it is only necessary to check if the sequence ID of the first - // chunk is duplicated. - // When we receive the initial message of a non-duplicated chunk message, we place it in the - // chunkMessageOngoing. Upon completion of sending this chunk message, if we receive other messages - // sent by this Producer, we will remove it from the chunkMessageOngoing. headersAndPayload.markReaderIndex(); MessageMetadata msgMetadata = Commands.parseMessageMetadata(headersAndPayload); headersAndPayload.resetReaderIndex(); - if (chunkMessageOngoing.containsKey(producerName)) { - //Deduplication requires the dependency on the ordered nature of messages to work accurately. - if (publishContext.isChunked() && chunkMessageOngoing.get(producerName) < msgMetadata.getChunkId()) { - chunkMessageOngoing.put(producerName, msgMetadata.getChunkId()); - return MessageDupStatus.NotDup; - } else { - chunkMessageOngoing.remove(producerName); - } - } long highestSequenceId = Math.max(publishContext.getHighestSequenceId(), sequenceId); if (producerName.startsWith(replicatorPrefix)) { // Message is coming from replication, we need to use the original producer name and sequence id @@ -361,12 +337,15 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade publishContext.setOriginalSequenceId(sequenceId); publishContext.setOriginalHighestSequenceId(highestSequenceId); } - + long chunkID = msgMetadata.hasChunkId() ? msgMetadata.getChunkId() : 0; // Synchronize the get() and subsequent put() on the map. This would only be relevant if the producer // disconnects and re-connects very quickly. At that point the call can be coming from a different thread synchronized (highestSequencedPushed) { Long lastSequenceIdPushed = highestSequencedPushed.get(producerName); - if (lastSequenceIdPushed != null && sequenceId <= lastSequenceIdPushed) { + // All chunks of a message use the same message metadata and sequence ID, + // so it's expected for sequenceId == lastSequenceIdPushed when the chunk ID > 0. + if (lastSequenceIdPushed != null && (chunkID > 0 ? sequenceId < lastSequenceIdPushed + : sequenceId <= lastSequenceIdPushed)) { if (log.isDebugEnabled()) { log.debug("[{}] Message identified as duplicated producer={} seq-id={} -- highest-seq-id={}", topic.getName(), producerName, sequenceId, lastSequenceIdPushed); @@ -379,7 +358,8 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade // lastSequenceIdPushed, then we cannot be sure whether the message is a dup or not // we should return an error to the producer for the latter case so that it can retry at a future time Long lastSequenceIdPersisted = highestSequencedPersisted.get(producerName); - if (lastSequenceIdPersisted != null && sequenceId <= lastSequenceIdPersisted) { + if (lastSequenceIdPersisted != null && (chunkID > 0 ? sequenceId < lastSequenceIdPersisted + : sequenceId <= lastSequenceIdPersisted)) { return MessageDupStatus.Dup; } else { return MessageDupStatus.Unknown; @@ -387,9 +367,6 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade } highestSequencedPushed.put(producerName, highestSequenceId); } - if (publishContext.isChunked()) { - chunkMessageOngoing.put(producerName, msgMetadata.getChunkId()); - } return MessageDupStatus.NotDup; } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java index cbc3ea68f0648..f726ef4d94676 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java @@ -246,14 +246,20 @@ public void testDuplicateForChunkMessage() throws Exception { msgIdGenerator.setAccessible(true); assertEquals(msg.getSequenceId() + 1, msgIdGenerator.get(partProducer)); - long sequenceID = (long) msgIdGenerator.get(partProducer) + 1024L; - String message2 = "b".repeat(messageSize * 1000); - partProducer.newMessage().value(message2).sequenceId(sequenceID).send(); + String message2 = "b".repeat(messageSize * 2); + partProducer.newMessage().value(message2).send(); Message msg2 = consumer.receive(5, TimeUnit.SECONDS); - assertNotNull(msg2); - assertTrue(msg2.getMessageId() instanceof ChunkMessageIdImpl); + assertFalse(msg2.getMessageId() instanceof ChunkMessageIdImpl); assertEquals(msg2.getValue(), message2); - assertEquals(msg2.getSequenceId(), sequenceID); + + long sequenceID = (long) msgIdGenerator.get(partProducer) + 1024L; + String message3 = "c".repeat(messageSize * 1000); + partProducer.newMessage().value(message3).sequenceId(sequenceID).send(); + Message msg3 = consumer.receive(5, TimeUnit.SECONDS); + assertNotNull(msg3); + assertTrue(msg3.getMessageId() instanceof ChunkMessageIdImpl); + assertEquals(msg3.getValue(), message3); + assertEquals(msg3.getSequenceId(), sequenceID); } private static void sendNonChunk(final PersistentTopic persistentTopic, From ab69604314cf7396354120b2592af14783fcfa34 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Fri, 18 Aug 2023 15:46:40 +0800 Subject: [PATCH 07/24] filter duplicated message at consumer side --- .../java/org/apache/pulsar/client/impl/ConsumerImpl.java | 9 +++++++++ 1 file changed, 9 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 a929fe9aa6bb2..90814e0c42b8f 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 @@ -1449,6 +1449,15 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m // discard message if chunk is out-of-order if (chunkedMsgCtx == null || chunkedMsgCtx.chunkedMsgBuffer == null || msgMetadata.getChunkId() != (chunkedMsgCtx.lastChunkedMessageId + 1)) { + //Filter duplicated chunks instead of discard it. + if (chunkedMsgCtx != null && msgMetadata.getChunkId() <= chunkedMsgCtx.lastChunkedMessageId) { + log.warn("[{}] Receive a repeated chunk messageId {}, last-chunk-id{}, chunkId = {}", + msgMetadata.getProducerName(), chunkedMsgCtx.lastChunkedMessageId, + msgId, msgMetadata.getChunkId()); + compressedPayload.release(); + increaseAvailablePermits(cnx); + return null; + } // means we lost the first chunk: should never happen log.info("Received unexpected chunk messageId {}, last-chunk-id{}, chunkId = {}", msgId, (chunkedMsgCtx != null ? chunkedMsgCtx.lastChunkedMessageId : null), msgMetadata.getChunkId()); From 3be7bf067fcf194754d2cea70bc4424e0cfb0a32 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Tue, 22 Aug 2023 21:53:46 +0800 Subject: [PATCH 08/24] add test --- .../impl/MessageChunkingSharedTest.java | 65 +++++++++++++++---- .../pulsar/client/impl/ConsumerImpl.java | 8 +-- 2 files changed, 57 insertions(+), 16 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java index f726ef4d94676..236b2577067a8 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java @@ -23,6 +23,7 @@ import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertTrue; import io.netty.buffer.ByteBuf; +import io.netty.buffer.Unpooled; import java.lang.reflect.Field; import java.time.Duration; import java.util.ArrayList; @@ -35,6 +36,7 @@ import java.util.concurrent.TimeUnit; import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; +import org.apache.pulsar.broker.service.Topic; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.ConsumerBuilder; @@ -43,11 +45,9 @@ import org.apache.pulsar.client.api.MessageListener; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.ProducerConsumerBase; -import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionType; -import org.apache.pulsar.common.allocator.PulsarByteBufAllocator; import org.apache.pulsar.common.api.proto.MessageMetadata; import org.apache.pulsar.common.protocol.Commands; import org.awaitility.Awaitility; @@ -219,7 +219,6 @@ public void testDuplicateForChunkMessage() throws Exception { restartBroker(); String topicName = "persistent://my-property/my-ns/testDuplicateForChunkMessage"; String producerName = "test-producer"; - pulsarClient = PulsarClient.builder().serviceUrl(pulsar.getBrokerServiceUrl()).build(); // consumer Consumer consumer = pulsarClient .newConsumer(Schema.STRING) @@ -262,6 +261,36 @@ public void testDuplicateForChunkMessage() throws Exception { assertEquals(msg3.getSequenceId(), sequenceID); } + @Test + public void testDeduplicateChunksInSingleChunkMessages() throws Exception { + this.conf.setBrokerDeduplicationEnabled(true); + restartBroker(); + String topicName = "persistent://my-property/my-ns/testDeduplicateChunksInSingleChunkMessage"; + String producerName = "test-producer"; + // consumer + Consumer consumer = pulsarClient + .newConsumer(Schema.STRING) + .subscriptionName("test-sub") + .topic(topicName) + .subscribe(); + final PersistentTopic persistentTopic = (PersistentTopic) pulsar.getBrokerService() + .getTopicIfExists(topicName).get().orElse(null); + assertNotNull(persistentTopic); + sendChunk(persistentTopic, "test-producer", 1, 0, 2); + sendChunk(persistentTopic, producerName, 1, 1, 2); + sendChunk(persistentTopic, producerName, 1, 1, 2); + + Message message = consumer.receive(10, TimeUnit.SECONDS); + assertEquals(message.getData().length, 2); + + sendChunk(persistentTopic, producerName, 2, 0, 3); + sendChunk(persistentTopic, producerName, 2, 1, 3); + sendChunk(persistentTopic, producerName, 2, 1, 3); + sendChunk(persistentTopic, producerName, 2, 2, 3); + message = consumer.receive(5, TimeUnit.SECONDS); + assertEquals(message.getData().length, 3); + } + private static void sendNonChunk(final PersistentTopic persistentTopic, final String producerName, final long sequenceId) { @@ -284,16 +313,28 @@ private static void sendChunk(final PersistentTopic persistentTopic, metadata.setTotalChunkMsgSize(numChunks); } final ByteBuf buf = Commands.serializeMetadataAndPayload(Commands.ChecksumType.Crc32c, metadata, - PulsarByteBufAllocator.DEFAULT.buffer(1)); - persistentTopic.publishMessage(buf, (e, ledgerId, entryId) -> { - String name = producerName + "-" + sequenceId; - if (chunkId != null) { - name += "-" + chunkId + "-" + numChunks; + Unpooled.wrappedBuffer("a".getBytes())); + persistentTopic.publishMessage(buf, new Topic.PublishContext() { + @Override + public String getProducerName() { + return producerName; + } + + public long getSequenceId() { + return sequenceId; } - if (e == null) { - log.info("Sent {} to ({}, {})", name, ledgerId, entryId); - } else { - log.error("Failed to send {}: {}", name, e.getMessage()); + + @Override + public void completed(Exception e, long ledgerId, long entryId) { + String name = producerName + "-" + sequenceId; + if (chunkId != null) { + name += "-" + chunkId + "-" + numChunks; + } + if (e == null) { + log.info("Sent {} to ({}, {})", name, ledgerId, entryId); + } else { + log.error("Failed to send {}: {}", name, e.getMessage()); + } } }); } 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 90814e0c42b8f..b45c692882d10 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 @@ -1449,11 +1449,11 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m // discard message if chunk is out-of-order if (chunkedMsgCtx == null || chunkedMsgCtx.chunkedMsgBuffer == null || msgMetadata.getChunkId() != (chunkedMsgCtx.lastChunkedMessageId + 1)) { - //Filter duplicated chunks instead of discard it. - if (chunkedMsgCtx != null && msgMetadata.getChunkId() <= chunkedMsgCtx.lastChunkedMessageId) { + // Filter duplicated chunks instead of discard it. + if (chunkedMsgCtx == null || msgMetadata.getChunkId() <= chunkedMsgCtx.lastChunkedMessageId) { log.warn("[{}] Receive a repeated chunk messageId {}, last-chunk-id{}, chunkId = {}", - msgMetadata.getProducerName(), chunkedMsgCtx.lastChunkedMessageId, - msgId, msgMetadata.getChunkId()); + msgMetadata.getProducerName(), chunkedMsgCtx == null ? null + : chunkedMsgCtx.lastChunkedMessageId, msgId, msgMetadata.getChunkId()); compressedPayload.release(); increaseAvailablePermits(cnx); return null; From 2b05a5c91cf70deaeb4ca731b097221ee977c026 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Wed, 23 Aug 2023 09:23:56 +0800 Subject: [PATCH 09/24] address some comment --- .../persistent/MessageDeduplication.java | 20 ++++++++++++------- .../impl/MessageChunkingSharedTest.java | 2 +- 2 files changed, 14 insertions(+), 8 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java index d51ab6b8eb07f..7857d2d14ecd5 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java @@ -323,21 +323,27 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade String producerName = publishContext.getProducerName(); long sequenceId = publishContext.getSequenceId(); - headersAndPayload.markReaderIndex(); - MessageMetadata msgMetadata = Commands.parseMessageMetadata(headersAndPayload); - headersAndPayload.resetReaderIndex(); long highestSequenceId = Math.max(publishContext.getHighestSequenceId(), sequenceId); if (producerName.startsWith(replicatorPrefix)) { // Message is coming from replication, we need to use the original producer name and sequence id // for the purpose of deduplication and not rely on the "replicator" name. - producerName = msgMetadata.getProducerName(); - sequenceId = msgMetadata.getSequenceId(); - highestSequenceId = Math.max(msgMetadata.getHighestSequenceId(), sequenceId); + int readerIndex = headersAndPayload.readerIndex(); + MessageMetadata md = Commands.parseMessageMetadata(headersAndPayload); + producerName = md.getProducerName(); + sequenceId = md.getSequenceId(); + highestSequenceId = Math.max(md.getHighestSequenceId(), sequenceId); publishContext.setOriginalProducerName(producerName); publishContext.setOriginalSequenceId(sequenceId); publishContext.setOriginalHighestSequenceId(highestSequenceId); + headersAndPayload.readerIndex(readerIndex); + } + long chunkID = 0; + if (publishContext.isChunked()) { + headersAndPayload.markReaderIndex(); + MessageMetadata msgMetadata = Commands.parseMessageMetadata(headersAndPayload); + headersAndPayload.resetReaderIndex(); + chunkID = msgMetadata.getChunkId(); } - long chunkID = msgMetadata.hasChunkId() ? msgMetadata.getChunkId() : 0; // Synchronize the get() and subsequent put() on the map. This would only be relevant if the producer // disconnects and re-connects very quickly. At that point the call can be coming from a different thread synchronized (highestSequencedPushed) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java index 236b2577067a8..0c88199270d85 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java @@ -276,7 +276,7 @@ public void testDeduplicateChunksInSingleChunkMessages() throws Exception { final PersistentTopic persistentTopic = (PersistentTopic) pulsar.getBrokerService() .getTopicIfExists(topicName).get().orElse(null); assertNotNull(persistentTopic); - sendChunk(persistentTopic, "test-producer", 1, 0, 2); + sendChunk(persistentTopic, producerName, 1, 0, 2); sendChunk(persistentTopic, producerName, 1, 1, 2); sendChunk(persistentTopic, producerName, 1, 1, 2); From 87a2513ca4fe839eccfbaa761c1aa81ffa64c30e Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Wed, 23 Aug 2023 10:02:43 +0800 Subject: [PATCH 10/24] address some comment --- .../pulsar/client/impl/ConsumerImpl.java | 24 +++++++++---------- 1 file changed, 11 insertions(+), 13 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 b45c692882d10..0181091d2526f 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 @@ -1454,25 +1454,23 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m log.warn("[{}] Receive a repeated chunk messageId {}, last-chunk-id{}, chunkId = {}", msgMetadata.getProducerName(), chunkedMsgCtx == null ? null : chunkedMsgCtx.lastChunkedMessageId, msgId, msgMetadata.getChunkId()); - compressedPayload.release(); - increaseAvailablePermits(cnx); - return null; - } - // means we lost the first chunk: should never happen - log.info("Received unexpected chunk messageId {}, last-chunk-id{}, chunkId = {}", msgId, - (chunkedMsgCtx != null ? chunkedMsgCtx.lastChunkedMessageId : null), msgMetadata.getChunkId()); - if (chunkedMsgCtx != null) { - if (chunkedMsgCtx.chunkedMsgBuffer != null) { - ReferenceCountUtil.safeRelease(chunkedMsgCtx.chunkedMsgBuffer); + } else { + // means we lost the first chunk: should never happen + log.info("Received unexpected chunk messageId {}, last-chunk-id{}, chunkId = {}", msgId, + (chunkedMsgCtx != null ? chunkedMsgCtx.lastChunkedMessageId : null), msgMetadata.getChunkId()); + if (chunkedMsgCtx != null) { + if (chunkedMsgCtx.chunkedMsgBuffer != null) { + ReferenceCountUtil.safeRelease(chunkedMsgCtx.chunkedMsgBuffer); + } + chunkedMsgCtx.recycle(); } - chunkedMsgCtx.recycle(); + chunkedMessagesMap.remove(msgMetadata.getUuid()); } - chunkedMessagesMap.remove(msgMetadata.getUuid()); compressedPayload.release(); increaseAvailablePermits(cnx); if (expireTimeOfIncompleteChunkedMessageMillis > 0 && System.currentTimeMillis() > (msgMetadata.getPublishTime() - + expireTimeOfIncompleteChunkedMessageMillis)) { + + expireTimeOfIncompleteChunkedMessageMillis)) { doAcknowledge(msgId, AckType.Individual, Collections.emptyMap(), null); } else { trackMessage(msgId); From 7af105397949e62e56f958b44a3e7c6023dae75e Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Wed, 23 Aug 2023 10:16:53 +0800 Subject: [PATCH 11/24] address some comment --- .../broker/service/persistent/MessageDeduplication.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java index 7857d2d14ecd5..7e0926d9a94f3 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java @@ -324,11 +324,12 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade String producerName = publishContext.getProducerName(); long sequenceId = publishContext.getSequenceId(); long highestSequenceId = Math.max(publishContext.getHighestSequenceId(), sequenceId); + MessageMetadata md = null; if (producerName.startsWith(replicatorPrefix)) { // Message is coming from replication, we need to use the original producer name and sequence id // for the purpose of deduplication and not rely on the "replicator" name. int readerIndex = headersAndPayload.readerIndex(); - MessageMetadata md = Commands.parseMessageMetadata(headersAndPayload); + md = Commands.parseMessageMetadata(headersAndPayload); producerName = md.getProducerName(); sequenceId = md.getSequenceId(); highestSequenceId = Math.max(md.getHighestSequenceId(), sequenceId); @@ -340,7 +341,7 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade long chunkID = 0; if (publishContext.isChunked()) { headersAndPayload.markReaderIndex(); - MessageMetadata msgMetadata = Commands.parseMessageMetadata(headersAndPayload); + MessageMetadata msgMetadata = (md == null) ? Commands.parseMessageMetadata(headersAndPayload) : md; headersAndPayload.resetReaderIndex(); chunkID = msgMetadata.getChunkId(); } From c4ef26f1483f4e64352e0e0df3fc0adc8fa98b85 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Wed, 23 Aug 2023 10:20:31 +0800 Subject: [PATCH 12/24] address some comment --- .../service/persistent/MessageDeduplication.java | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java index 7e0926d9a94f3..9c2c47caaee18 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java @@ -340,10 +340,12 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade } long chunkID = 0; if (publishContext.isChunked()) { - headersAndPayload.markReaderIndex(); - MessageMetadata msgMetadata = (md == null) ? Commands.parseMessageMetadata(headersAndPayload) : md; - headersAndPayload.resetReaderIndex(); - chunkID = msgMetadata.getChunkId(); + if (md == null) { + headersAndPayload.markReaderIndex(); + md = Commands.parseMessageMetadata(headersAndPayload); + headersAndPayload.resetReaderIndex(); + } + chunkID = md.getChunkId(); } // Synchronize the get() and subsequent put() on the map. This would only be relevant if the producer // disconnects and re-connects very quickly. At that point the call can be coming from a different thread From fb1a55c391e43d0b04bf4ecaa0202f70153db74a Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Thu, 24 Aug 2023 08:29:31 +0800 Subject: [PATCH 13/24] Revert "address some comment" This reverts commit 87a2513ca4fe839eccfbaa761c1aa81ffa64c30e. --- .../pulsar/client/impl/ConsumerImpl.java | 24 ++++++++++--------- 1 file changed, 13 insertions(+), 11 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 0181091d2526f..b45c692882d10 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 @@ -1454,23 +1454,25 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m log.warn("[{}] Receive a repeated chunk messageId {}, last-chunk-id{}, chunkId = {}", msgMetadata.getProducerName(), chunkedMsgCtx == null ? null : chunkedMsgCtx.lastChunkedMessageId, msgId, msgMetadata.getChunkId()); - } else { - // means we lost the first chunk: should never happen - log.info("Received unexpected chunk messageId {}, last-chunk-id{}, chunkId = {}", msgId, - (chunkedMsgCtx != null ? chunkedMsgCtx.lastChunkedMessageId : null), msgMetadata.getChunkId()); - if (chunkedMsgCtx != null) { - if (chunkedMsgCtx.chunkedMsgBuffer != null) { - ReferenceCountUtil.safeRelease(chunkedMsgCtx.chunkedMsgBuffer); - } - chunkedMsgCtx.recycle(); + compressedPayload.release(); + increaseAvailablePermits(cnx); + return null; + } + // means we lost the first chunk: should never happen + log.info("Received unexpected chunk messageId {}, last-chunk-id{}, chunkId = {}", msgId, + (chunkedMsgCtx != null ? chunkedMsgCtx.lastChunkedMessageId : null), msgMetadata.getChunkId()); + if (chunkedMsgCtx != null) { + if (chunkedMsgCtx.chunkedMsgBuffer != null) { + ReferenceCountUtil.safeRelease(chunkedMsgCtx.chunkedMsgBuffer); } - chunkedMessagesMap.remove(msgMetadata.getUuid()); + chunkedMsgCtx.recycle(); } + chunkedMessagesMap.remove(msgMetadata.getUuid()); compressedPayload.release(); increaseAvailablePermits(cnx); if (expireTimeOfIncompleteChunkedMessageMillis > 0 && System.currentTimeMillis() > (msgMetadata.getPublishTime() - + expireTimeOfIncompleteChunkedMessageMillis)) { + + expireTimeOfIncompleteChunkedMessageMillis)) { doAcknowledge(msgId, AckType.Individual, Collections.emptyMap(), null); } else { trackMessage(msgId); From d9af22f09de99070ab424a72cbfb0b07277efaa3 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Thu, 24 Aug 2023 22:57:19 +0800 Subject: [PATCH 14/24] ack duplicated chunk --- .../java/org/apache/pulsar/client/impl/ConsumerImpl.java | 9 +++++++++ 1 file changed, 9 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 b45c692882d10..56be57cc94091 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 @@ -35,6 +35,7 @@ import java.io.IOException; import java.nio.ByteBuffer; import java.util.ArrayList; +import java.util.Arrays; import java.util.BitSet; import java.util.Collections; import java.util.HashMap; @@ -1456,6 +1457,14 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m : chunkedMsgCtx.lastChunkedMessageId, msgId, msgMetadata.getChunkId()); compressedPayload.release(); increaseAvailablePermits(cnx); + if (chunkedMsgCtx != null) { + boolean repeatedlyReceived = Arrays.stream(chunkedMsgCtx.chunkedMessageIds) + .anyMatch(messageId1 -> messageId1.ledgerId == messageId.getLedgerId() + && messageId1.entryId == messageId.getEntryId()); + if (!repeatedlyReceived) { + doAcknowledge(msgId, AckType.Individual, Collections.emptyMap(), null); + } + } return null; } // means we lost the first chunk: should never happen From 86557230a36b165f80c7460e0c62b71185b42c78 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Fri, 25 Aug 2023 09:13:46 +0800 Subject: [PATCH 15/24] rollback the logic of `chunkedMsgCtx == null` and add note. --- .../persistent/MessageDeduplication.java | 4 ++++ .../pulsar/client/impl/ConsumerImpl.java | 23 ++++++++++--------- 2 files changed, 16 insertions(+), 11 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java index 9c2c47caaee18..78276d65170d2 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java @@ -353,6 +353,10 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade Long lastSequenceIdPushed = highestSequencedPushed.get(producerName); // All chunks of a message use the same message metadata and sequence ID, // so it's expected for sequenceId == lastSequenceIdPushed when the chunk ID > 0. + // "chunkID == 0" means that the message is the first one of the chunk list. + // We check the sequence ID of the first chunk as the same as the common messages. + // Todo: Add the last chunkID map (like `highestSequencedPushed` and `highestSequencedPersisted`) to check + // the duplication in the chunk list of a chunk message. if (lastSequenceIdPushed != null && (chunkID > 0 ? sequenceId < lastSequenceIdPushed : sequenceId <= lastSequenceIdPushed)) { if (log.isDebugEnabled()) { 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 56be57cc94091..aaf5d1b6ff037 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 @@ -1450,20 +1450,21 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m // discard message if chunk is out-of-order if (chunkedMsgCtx == null || chunkedMsgCtx.chunkedMsgBuffer == null || msgMetadata.getChunkId() != (chunkedMsgCtx.lastChunkedMessageId + 1)) { - // Filter duplicated chunks instead of discard it. - if (chunkedMsgCtx == null || msgMetadata.getChunkId() <= chunkedMsgCtx.lastChunkedMessageId) { + // Filter duplicated chunks instead of discard it. (Only do this when exist duplication in a chunk message) + // For example: + // Chunk-1 sequence ID: 0, chunk ID: 0 + // Chunk-2 sequence ID: 0, chunk ID: 0 + // Chunk-3 sequence ID: 0, chunk ID: 1 + if (chunkedMsgCtx != null && msgMetadata.getChunkId() <= chunkedMsgCtx.lastChunkedMessageId) { log.warn("[{}] Receive a repeated chunk messageId {}, last-chunk-id{}, chunkId = {}", - msgMetadata.getProducerName(), chunkedMsgCtx == null ? null - : chunkedMsgCtx.lastChunkedMessageId, msgId, msgMetadata.getChunkId()); + msgMetadata.getProducerName(), chunkedMsgCtx.lastChunkedMessageId, msgId, msgMetadata.getChunkId()); compressedPayload.release(); increaseAvailablePermits(cnx); - if (chunkedMsgCtx != null) { - boolean repeatedlyReceived = Arrays.stream(chunkedMsgCtx.chunkedMessageIds) - .anyMatch(messageId1 -> messageId1.ledgerId == messageId.getLedgerId() - && messageId1.entryId == messageId.getEntryId()); - if (!repeatedlyReceived) { - doAcknowledge(msgId, AckType.Individual, Collections.emptyMap(), null); - } + boolean repeatedlyReceived = Arrays.stream(chunkedMsgCtx.chunkedMessageIds) + .anyMatch(messageId1 -> messageId1.ledgerId == messageId.getLedgerId() + && messageId1.entryId == messageId.getEntryId()); + if (!repeatedlyReceived) { + doAcknowledge(msgId, AckType.Individual, Collections.emptyMap(), null); } return null; } From d7f9290c48a49ebff0aee2de090542970c56ea9e Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Fri, 25 Aug 2023 10:29:04 +0800 Subject: [PATCH 16/24] fix test --- .../apache/pulsar/client/impl/MessageChunkingSharedTest.java | 5 +++++ .../java/org/apache/pulsar/client/impl/ConsumerImpl.java | 2 +- 2 files changed, 6 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java index 0c88199270d85..2569a59369826 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java @@ -315,6 +315,11 @@ private static void sendChunk(final PersistentTopic persistentTopic, final ByteBuf buf = Commands.serializeMetadataAndPayload(Commands.ChecksumType.Crc32c, metadata, Unpooled.wrappedBuffer("a".getBytes())); persistentTopic.publishMessage(buf, new Topic.PublishContext() { + @Override + public boolean isChunked() { + return chunkId != null; + } + @Override public String getProducerName() { return producerName; 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 aaf5d1b6ff037..2fa64c53fc06d 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 @@ -1461,7 +1461,7 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m compressedPayload.release(); increaseAvailablePermits(cnx); boolean repeatedlyReceived = Arrays.stream(chunkedMsgCtx.chunkedMessageIds) - .anyMatch(messageId1 -> messageId1.ledgerId == messageId.getLedgerId() + .anyMatch(messageId1 -> messageId1 != null && messageId1.ledgerId == messageId.getLedgerId() && messageId1.entryId == messageId.getEntryId()); if (!repeatedlyReceived) { doAcknowledge(msgId, AckType.Individual, Collections.emptyMap(), null); From f282ac0c462b731a35f703a746d13c798b4c3eb2 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Sat, 26 Aug 2023 09:33:10 +0800 Subject: [PATCH 17/24] optimize uuid --- .../main/java/org/apache/pulsar/client/impl/ProducerImpl.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java index 267b06649d719..cd840740be50e 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java @@ -554,7 +554,8 @@ public void sendAsync(Message message, SendCallback callback) { // Update the message metadata before computing the payload chunk size // to avoid a large message cannot be split into chunks. final long sequenceId = updateMessageMetadataSequenceId(msgMetadata); - String uuid = totalChunks > 1 ? String.format("%s-%d", producerName, sequenceId) : null; + String uuid = totalChunks > 1 ? String.format("%s-%d-%d", producerName, sequenceId, + System.currentTimeMillis()) : null; serializeAndSendMessage(msg, payload, sequenceId, uuid, chunkId, totalChunks, readStartIndex, payloadChunkSize, compressedPayload, compressed, From 1f8263d630942646f4ba53663b290727ef7b12c2 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Mon, 28 Aug 2023 08:56:57 +0800 Subject: [PATCH 18/24] optimize uuid and checkstyle --- .../main/java/org/apache/pulsar/client/impl/ConsumerImpl.java | 3 ++- .../main/java/org/apache/pulsar/client/impl/ProducerImpl.java | 3 ++- 2 files changed, 4 insertions(+), 2 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 2fa64c53fc06d..df5bf9e96a74c 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 @@ -1457,7 +1457,8 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m // Chunk-3 sequence ID: 0, chunk ID: 1 if (chunkedMsgCtx != null && msgMetadata.getChunkId() <= chunkedMsgCtx.lastChunkedMessageId) { log.warn("[{}] Receive a repeated chunk messageId {}, last-chunk-id{}, chunkId = {}", - msgMetadata.getProducerName(), chunkedMsgCtx.lastChunkedMessageId, msgId, msgMetadata.getChunkId()); + msgMetadata.getProducerName(), chunkedMsgCtx.lastChunkedMessageId, + msgId, msgMetadata.getChunkId()); compressedPayload.release(); increaseAvailablePermits(cnx); boolean repeatedlyReceived = Arrays.stream(chunkedMsgCtx.chunkedMessageIds) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java index cd840740be50e..d74d18be18338 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java @@ -531,6 +531,7 @@ public void sendAsync(Message message, SendCallback callback) { ? msg.getMessageBuilder().getOrderingKey() : null; // msg.messageId will be reset if previous message chunk is sent successfully. final MessageId messageId = msg.getMessageId(); + long timestamp = System.currentTimeMillis(); for (int chunkId = 0; chunkId < totalChunks; chunkId++) { // Need to reset the schemaVersion, because the schemaVersion is based on a ByteBuf object in // `MessageMetadata`, if we want to re-serialize the `SEND` command using a same `MessageMetadata`, @@ -555,7 +556,7 @@ public void sendAsync(Message message, SendCallback callback) { // to avoid a large message cannot be split into chunks. final long sequenceId = updateMessageMetadataSequenceId(msgMetadata); String uuid = totalChunks > 1 ? String.format("%s-%d-%d", producerName, sequenceId, - System.currentTimeMillis()) : null; + timestamp) : null; serializeAndSendMessage(msg, payload, sequenceId, uuid, chunkId, totalChunks, readStartIndex, payloadChunkSize, compressedPayload, compressed, From 4f789ee73e506b1e6e4653a9037636cc8ae85828 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Mon, 28 Aug 2023 10:20:46 +0800 Subject: [PATCH 19/24] check the last chunk ID --- .../broker/service/persistent/MessageDeduplication.java | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java index 78276d65170d2..baf2c64f2528d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java @@ -338,7 +338,8 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade publishContext.setOriginalHighestSequenceId(highestSequenceId); headersAndPayload.readerIndex(readerIndex); } - long chunkID = 0; + long chunkID = -1; + long totalChunk = -1; if (publishContext.isChunked()) { if (md == null) { headersAndPayload.markReaderIndex(); @@ -346,6 +347,7 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade headersAndPayload.resetReaderIndex(); } chunkID = md.getChunkId(); + totalChunk = md.getNumChunksFromMsg(); } // Synchronize the get() and subsequent put() on the map. This would only be relevant if the producer // disconnects and re-connects very quickly. At that point the call can be coming from a different thread @@ -357,7 +359,8 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade // We check the sequence ID of the first chunk as the same as the common messages. // Todo: Add the last chunkID map (like `highestSequencedPushed` and `highestSequencedPersisted`) to check // the duplication in the chunk list of a chunk message. - if (lastSequenceIdPushed != null && (chunkID > 0 ? sequenceId < lastSequenceIdPushed + if (lastSequenceIdPushed != null + && (chunkID >= 0 && chunkID != totalChunk - 1 ? sequenceId < lastSequenceIdPushed : sequenceId <= lastSequenceIdPushed)) { if (log.isDebugEnabled()) { log.debug("[{}] Message identified as duplicated producer={} seq-id={} -- highest-seq-id={}", From 2aaac4596da5d6fedde6c5ad1830b951c95b8ccc Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Wed, 30 Aug 2023 16:06:52 +0800 Subject: [PATCH 20/24] only store the sequence ID of the last chunk ID --- .../persistent/MessageDeduplication.java | 36 ++++++++++------- .../impl/MessageChunkingSharedTest.java | 39 +++++++++++++++++-- .../pulsar/client/impl/ConsumerImpl.java | 14 ++++--- .../pulsar/client/impl/ProducerImpl.java | 4 +- 4 files changed, 67 insertions(+), 26 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java index 6e16a0f815b62..238dc740509b1 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java @@ -56,6 +56,8 @@ public class MessageDeduplication { private final ManagedLedger managedLedger; private ManagedCursor managedCursor; + private static final String IS_LAST_CHUNK = "isLastChunk"; + enum Status { // Deduplication is initialized @@ -346,26 +348,24 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade long totalChunk = -1; if (publishContext.isChunked()) { if (md == null) { - headersAndPayload.markReaderIndex(); + int readerIndex = headersAndPayload.readerIndex(); md = Commands.parseMessageMetadata(headersAndPayload); - headersAndPayload.resetReaderIndex(); + headersAndPayload.readerIndex(readerIndex); } chunkID = md.getChunkId(); totalChunk = md.getNumChunksFromMsg(); } + // All chunks of a message use the same message metadata and sequence ID, + // so we only need to check the sequence ID for the last chunk in a chunk message. + if (chunkID != -1 && chunkID != totalChunk - 1) { + publishContext.setProperty(IS_LAST_CHUNK, Boolean.FALSE); + return MessageDupStatus.NotDup; + } // Synchronize the get() and subsequent put() on the map. This would only be relevant if the producer // disconnects and re-connects very quickly. At that point the call can be coming from a different thread synchronized (highestSequencedPushed) { Long lastSequenceIdPushed = highestSequencedPushed.get(producerName); - // All chunks of a message use the same message metadata and sequence ID, - // so it's expected for sequenceId == lastSequenceIdPushed when the chunk ID > 0. - // "chunkID == 0" means that the message is the first one of the chunk list. - // We check the sequence ID of the first chunk as the same as the common messages. - // Todo: Add the last chunkID map (like `highestSequencedPushed` and `highestSequencedPersisted`) to check - // the duplication in the chunk list of a chunk message. - if (lastSequenceIdPushed != null - && (chunkID >= 0 && chunkID != totalChunk - 1 ? sequenceId < lastSequenceIdPushed - : sequenceId <= lastSequenceIdPushed)) { + if (lastSequenceIdPushed != null && sequenceId <= lastSequenceIdPushed) { if (log.isDebugEnabled()) { log.debug("[{}] Message identified as duplicated producer={} seq-id={} -- highest-seq-id={}", topic.getName(), producerName, sequenceId, lastSequenceIdPushed); @@ -378,8 +378,7 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade // lastSequenceIdPushed, then we cannot be sure whether the message is a dup or not // we should return an error to the producer for the latter case so that it can retry at a future time Long lastSequenceIdPersisted = highestSequencedPersisted.get(producerName); - if (lastSequenceIdPersisted != null && (chunkID > 0 ? sequenceId < lastSequenceIdPersisted - : sequenceId <= lastSequenceIdPersisted)) { + if (lastSequenceIdPersisted != null && sequenceId <= lastSequenceIdPersisted) { return MessageDupStatus.Dup; } else { return MessageDupStatus.Unknown; @@ -387,6 +386,11 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade } highestSequencedPushed.put(producerName, highestSequenceId); } + // Only put sequence ID into highestSequencedPushed and + // highestSequencedPersisted until receive and persistent the last chunk. + if (chunkID != -1 && chunkID == totalChunk - 1) { + publishContext.setProperty(IS_LAST_CHUNK, Boolean.TRUE); + } return MessageDupStatus.NotDup; } @@ -407,8 +411,10 @@ public void recordMessagePersisted(PublishContext publishContext, PositionImpl p sequenceId = publishContext.getOriginalSequenceId(); highestSequenceId = publishContext.getOriginalHighestSequenceId(); } - - highestSequencedPersisted.put(producerName, Math.max(highestSequenceId, sequenceId)); + Boolean isLastChunk = (Boolean) publishContext.getProperty(IS_LAST_CHUNK); + if (isLastChunk == null || isLastChunk) { + highestSequencedPersisted.put(producerName, Math.max(highestSequenceId, sequenceId)); + } if (++snapshotCounter >= snapshotInterval) { snapshotCounter = 0; takeSnapshot(position); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java index 2569a59369826..e029dc4b42609 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java @@ -21,6 +21,7 @@ import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertNotNull; +import static org.testng.Assert.assertNull; import static org.testng.Assert.assertTrue; import io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; @@ -213,19 +214,51 @@ private static String createChunkedMessage(int numChunks) { return Schema.STRING.decode(payload); } + @Test + public void testSendChunkMessageWithSameSequenceID() throws Exception { + this.conf.setBrokerDeduplicationEnabled(true); + restartBroker(); + String topicName = "persistent://my-property/my-ns/testSendChunkMessageWithSameSequenceID"; + String producerName = "test-producer"; + @Cleanup + Consumer consumer = pulsarClient + .newConsumer(Schema.STRING) + .subscriptionName("test-sub") + .topic(topicName) + .subscribe(); + @Cleanup + Producer producer = pulsarClient + .newProducer(Schema.STRING) + .producerName(producerName) + .topic(topicName) + .enableChunking(true) + .enableBatching(false) + .create(); + int messageSize = 6000; // payload size in KB + String message = "a".repeat(messageSize * 1000); + producer.newMessage().value(message).sequenceId(10).send(); + Message msg = consumer.receive(5, TimeUnit.SECONDS); + assertNotNull(msg); + assertTrue(msg.getMessageId() instanceof ChunkMessageIdImpl); + assertEquals(msg.getValue(), message); + producer.newMessage().value(message).sequenceId(10).send(); + msg = consumer.receive(3, TimeUnit.SECONDS); + assertNull(msg); + } + @Test public void testDuplicateForChunkMessage() throws Exception { this.conf.setBrokerDeduplicationEnabled(true); restartBroker(); String topicName = "persistent://my-property/my-ns/testDuplicateForChunkMessage"; String producerName = "test-producer"; - // consumer + @Cleanup Consumer consumer = pulsarClient .newConsumer(Schema.STRING) .subscriptionName("test-sub") .topic(topicName) .subscribe(); - // producer + @Cleanup Producer partProducer = pulsarClient .newProducer(Schema.STRING) .producerName(producerName) @@ -267,7 +300,7 @@ public void testDeduplicateChunksInSingleChunkMessages() throws Exception { restartBroker(); String topicName = "persistent://my-property/my-ns/testDeduplicateChunksInSingleChunkMessage"; String producerName = "test-producer"; - // consumer + @Cleanup Consumer consumer = pulsarClient .newConsumer(Schema.STRING) .subscriptionName("test-sub") 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 81ee1b14d5977..2bacb40424692 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 @@ -1460,11 +1460,15 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m // discard message if chunk is out-of-order if (chunkedMsgCtx == null || chunkedMsgCtx.chunkedMsgBuffer == null || msgMetadata.getChunkId() != (chunkedMsgCtx.lastChunkedMessageId + 1)) { - // Filter duplicated chunks instead of discard it. (Only do this when exist duplication in a chunk message) - // For example: - // Chunk-1 sequence ID: 0, chunk ID: 0 - // Chunk-2 sequence ID: 0, chunk ID: 0 - // Chunk-3 sequence ID: 0, chunk ID: 1 +// Filter and ack duplicated chunks instead of discard ctx. +// 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: 2, 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 +// Chunk-6 sequence ID: 0, chunk ID: 3, msgID: 1:6 +// We should filter and ack chunk-4 and chunk-5. if (chunkedMsgCtx != null && msgMetadata.getChunkId() <= chunkedMsgCtx.lastChunkedMessageId) { log.warn("[{}] Receive a repeated chunk messageId {}, last-chunk-id{}, chunkId = {}", msgMetadata.getProducerName(), chunkedMsgCtx.lastChunkedMessageId, diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java index 90e9faed7809c..9bd1421d0355b 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java @@ -531,7 +531,6 @@ public void sendAsync(Message message, SendCallback callback) { ? msg.getMessageBuilder().getOrderingKey() : null; // msg.messageId will be reset if previous message chunk is sent successfully. final MessageId messageId = msg.getMessageId(); - long timestamp = System.currentTimeMillis(); for (int chunkId = 0; chunkId < totalChunks; chunkId++) { // Need to reset the schemaVersion, because the schemaVersion is based on a ByteBuf object in // `MessageMetadata`, if we want to re-serialize the `SEND` command using a same `MessageMetadata`, @@ -555,8 +554,7 @@ public void sendAsync(Message message, SendCallback callback) { // Update the message metadata before computing the payload chunk size // to avoid a large message cannot be split into chunks. final long sequenceId = updateMessageMetadataSequenceId(msgMetadata); - String uuid = totalChunks > 1 ? String.format("%s-%d-%d", producerName, sequenceId, - timestamp) : null; + String uuid = totalChunks > 1 ? String.format("%s-%d", producerName, sequenceId) : null; serializeAndSendMessage(msg, payload, sequenceId, uuid, chunkId, totalChunks, readStartIndex, payloadChunkSize, compressedPayload, compressed, From 8579c3707a4860209d82b5d90aad5040d29edcab Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Wed, 30 Aug 2023 16:26:20 +0800 Subject: [PATCH 21/24] Avoid flaky testing --- .../pulsar/client/impl/MessageChunkingSharedTest.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java index e029dc4b42609..85ef717cdf91f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java @@ -237,7 +237,7 @@ public void testSendChunkMessageWithSameSequenceID() throws Exception { int messageSize = 6000; // payload size in KB String message = "a".repeat(messageSize * 1000); producer.newMessage().value(message).sequenceId(10).send(); - Message msg = consumer.receive(5, TimeUnit.SECONDS); + Message msg = consumer.receive(10, TimeUnit.SECONDS); assertNotNull(msg); assertTrue(msg.getMessageId() instanceof ChunkMessageIdImpl); assertEquals(msg.getValue(), message); @@ -313,14 +313,14 @@ public void testDeduplicateChunksInSingleChunkMessages() throws Exception { sendChunk(persistentTopic, producerName, 1, 1, 2); sendChunk(persistentTopic, producerName, 1, 1, 2); - Message message = consumer.receive(10, TimeUnit.SECONDS); + Message message = consumer.receive(15, TimeUnit.SECONDS); assertEquals(message.getData().length, 2); sendChunk(persistentTopic, producerName, 2, 0, 3); sendChunk(persistentTopic, producerName, 2, 1, 3); sendChunk(persistentTopic, producerName, 2, 1, 3); sendChunk(persistentTopic, producerName, 2, 2, 3); - message = consumer.receive(5, TimeUnit.SECONDS); + message = consumer.receive(20, TimeUnit.SECONDS); assertEquals(message.getData().length, 3); } From 3e2267c8f9b29cdc124849564d56e71d938d5267 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Wed, 30 Aug 2023 19:39:01 +0800 Subject: [PATCH 22/24] Make comments indent aligned. --- .../pulsar/client/impl/ConsumerImpl.java | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 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 2bacb40424692..58b97bf158959 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 @@ -1460,15 +1460,15 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m // discard message if chunk is out-of-order if (chunkedMsgCtx == null || chunkedMsgCtx.chunkedMsgBuffer == null || msgMetadata.getChunkId() != (chunkedMsgCtx.lastChunkedMessageId + 1)) { -// Filter and ack duplicated chunks instead of discard ctx. -// 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: 2, 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 -// Chunk-6 sequence ID: 0, chunk ID: 3, msgID: 1:6 -// We should filter and ack chunk-4 and chunk-5. + // Filter and ack duplicated chunks instead of discard ctx. + // 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: 2, 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 + // Chunk-6 sequence ID: 0, chunk ID: 3, msgID: 1:6 + // We should filter and ack chunk-4 and chunk-5. if (chunkedMsgCtx != null && msgMetadata.getChunkId() <= chunkedMsgCtx.lastChunkedMessageId) { log.warn("[{}] Receive a repeated chunk messageId {}, last-chunk-id{}, chunkId = {}", msgMetadata.getProducerName(), chunkedMsgCtx.lastChunkedMessageId, From 7c813e9a8ccae9df857f086b47ca71bcf75175aa Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Wed, 30 Aug 2023 20:17:20 +0800 Subject: [PATCH 23/24] Add a new test class --- .../MessageChunkingDeduplicationTest.java | 162 ++++++++++++++++++ .../impl/MessageChunkingSharedTest.java | 114 +----------- 2 files changed, 163 insertions(+), 113 deletions(-) create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingDeduplicationTest.java diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingDeduplicationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingDeduplicationTest.java new file mode 100644 index 0000000000000..24aa2356017eb --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingDeduplicationTest.java @@ -0,0 +1,162 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.client.impl; + +import static org.apache.pulsar.client.impl.MessageChunkingSharedTest.sendChunk; +import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertFalse; +import static org.testng.Assert.assertNotNull; +import static org.testng.Assert.assertNull; +import static org.testng.Assert.assertTrue; +import java.lang.reflect.Field; +import java.util.concurrent.TimeUnit; +import lombok.Cleanup; +import lombok.extern.slf4j.Slf4j; +import org.apache.pulsar.broker.service.persistent.PersistentTopic; +import org.apache.pulsar.client.api.Consumer; +import org.apache.pulsar.client.api.Message; +import org.apache.pulsar.client.api.Producer; +import org.apache.pulsar.client.api.ProducerConsumerBase; +import org.apache.pulsar.client.api.Schema; +import org.testng.annotations.AfterClass; +import org.testng.annotations.BeforeClass; +import org.testng.annotations.Test; + +@Slf4j +@Test(groups = "broker-impl") +public class MessageChunkingDeduplicationTest extends ProducerConsumerBase { + + @BeforeClass + @Override + protected void setup() throws Exception { + this.conf.setBrokerDeduplicationEnabled(true); + super.internalSetup(); + super.producerBaseSetup(); + } + + @AfterClass(alwaysRun = true) + @Override + protected void cleanup() throws Exception { + super.internalCleanup(); + } + + @Test + public void testSendChunkMessageWithSameSequenceID() throws Exception { + String topicName = "persistent://my-property/my-ns/testSendChunkMessageWithSameSequenceID"; + String producerName = "test-producer"; + @Cleanup + Consumer consumer = pulsarClient + .newConsumer(Schema.STRING) + .subscriptionName("test-sub") + .topic(topicName) + .subscribe(); + @Cleanup + Producer producer = pulsarClient + .newProducer(Schema.STRING) + .producerName(producerName) + .topic(topicName) + .enableChunking(true) + .enableBatching(false) + .create(); + int messageSize = 6000; // payload size in KB + String message = "a".repeat(messageSize * 1000); + producer.newMessage().value(message).sequenceId(10).send(); + Message msg = consumer.receive(10, TimeUnit.SECONDS); + assertNotNull(msg); + assertTrue(msg.getMessageId() instanceof ChunkMessageIdImpl); + assertEquals(msg.getValue(), message); + producer.newMessage().value(message).sequenceId(10).send(); + msg = consumer.receive(3, TimeUnit.SECONDS); + assertNull(msg); + } + + @Test + public void testDuplicateForChunkMessage() throws Exception { + String topicName = "persistent://my-property/my-ns/testDuplicateForChunkMessage"; + String producerName = "test-producer"; + @Cleanup + Consumer consumer = pulsarClient + .newConsumer(Schema.STRING) + .subscriptionName("test-sub") + .topic(topicName) + .subscribe(); + @Cleanup + Producer partProducer = pulsarClient + .newProducer(Schema.STRING) + .producerName(producerName) + .topic(topicName) + .enableChunking(true) + .enableBatching(false) + .create(); + int messageSize = 6000; // payload size in KB + String message = "a".repeat(messageSize * 1000); + partProducer.newMessage().value(message).send(); + Message msg = consumer.receive(5, TimeUnit.SECONDS); + assertNotNull(msg); + assertTrue(msg.getMessageId() instanceof ChunkMessageIdImpl); + assertEquals(msg.getValue(), message); + + Field msgIdGenerator = ProducerImpl.class.getDeclaredField("msgIdGenerator"); + msgIdGenerator.setAccessible(true); + assertEquals(msg.getSequenceId() + 1, msgIdGenerator.get(partProducer)); + + String message2 = "b".repeat(messageSize * 2); + partProducer.newMessage().value(message2).send(); + Message msg2 = consumer.receive(5, TimeUnit.SECONDS); + assertFalse(msg2.getMessageId() instanceof ChunkMessageIdImpl); + assertEquals(msg2.getValue(), message2); + + long sequenceID = (long) msgIdGenerator.get(partProducer) + 1024L; + String message3 = "c".repeat(messageSize * 1000); + partProducer.newMessage().value(message3).sequenceId(sequenceID).send(); + Message msg3 = consumer.receive(5, TimeUnit.SECONDS); + assertNotNull(msg3); + assertTrue(msg3.getMessageId() instanceof ChunkMessageIdImpl); + assertEquals(msg3.getValue(), message3); + assertEquals(msg3.getSequenceId(), sequenceID); + } + + @Test + public void testDeduplicateChunksInSingleChunkMessages() throws Exception { + String topicName = "persistent://my-property/my-ns/testDeduplicateChunksInSingleChunkMessage"; + String producerName = "test-producer"; + @Cleanup + Consumer consumer = pulsarClient + .newConsumer(Schema.STRING) + .subscriptionName("test-sub") + .topic(topicName) + .subscribe(); + final PersistentTopic persistentTopic = (PersistentTopic) pulsar.getBrokerService() + .getTopicIfExists(topicName).get().orElse(null); + assertNotNull(persistentTopic); + sendChunk(persistentTopic, producerName, 1, 0, 2); + sendChunk(persistentTopic, producerName, 1, 1, 2); + sendChunk(persistentTopic, producerName, 1, 1, 2); + + Message message = consumer.receive(15, TimeUnit.SECONDS); + assertEquals(message.getData().length, 2); + + sendChunk(persistentTopic, producerName, 2, 0, 3); + sendChunk(persistentTopic, producerName, 2, 1, 3); + sendChunk(persistentTopic, producerName, 2, 1, 3); + sendChunk(persistentTopic, producerName, 2, 2, 3); + message = consumer.receive(20, TimeUnit.SECONDS); + assertEquals(message.getData().length, 3); + } +} diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java index 85ef717cdf91f..3d24d3746d66a 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingSharedTest.java @@ -21,11 +21,9 @@ import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertNotNull; -import static org.testng.Assert.assertNull; import static org.testng.Assert.assertTrue; import io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; -import java.lang.reflect.Field; import java.time.Duration; import java.util.ArrayList; import java.util.Arrays; @@ -214,123 +212,13 @@ private static String createChunkedMessage(int numChunks) { return Schema.STRING.decode(payload); } - @Test - public void testSendChunkMessageWithSameSequenceID() throws Exception { - this.conf.setBrokerDeduplicationEnabled(true); - restartBroker(); - String topicName = "persistent://my-property/my-ns/testSendChunkMessageWithSameSequenceID"; - String producerName = "test-producer"; - @Cleanup - Consumer consumer = pulsarClient - .newConsumer(Schema.STRING) - .subscriptionName("test-sub") - .topic(topicName) - .subscribe(); - @Cleanup - Producer producer = pulsarClient - .newProducer(Schema.STRING) - .producerName(producerName) - .topic(topicName) - .enableChunking(true) - .enableBatching(false) - .create(); - int messageSize = 6000; // payload size in KB - String message = "a".repeat(messageSize * 1000); - producer.newMessage().value(message).sequenceId(10).send(); - Message msg = consumer.receive(10, TimeUnit.SECONDS); - assertNotNull(msg); - assertTrue(msg.getMessageId() instanceof ChunkMessageIdImpl); - assertEquals(msg.getValue(), message); - producer.newMessage().value(message).sequenceId(10).send(); - msg = consumer.receive(3, TimeUnit.SECONDS); - assertNull(msg); - } - - @Test - public void testDuplicateForChunkMessage() throws Exception { - this.conf.setBrokerDeduplicationEnabled(true); - restartBroker(); - String topicName = "persistent://my-property/my-ns/testDuplicateForChunkMessage"; - String producerName = "test-producer"; - @Cleanup - Consumer consumer = pulsarClient - .newConsumer(Schema.STRING) - .subscriptionName("test-sub") - .topic(topicName) - .subscribe(); - @Cleanup - Producer partProducer = pulsarClient - .newProducer(Schema.STRING) - .producerName(producerName) - .topic(topicName) - .enableChunking(true) - .enableBatching(false) - .create(); - int messageSize = 6000; // payload size in KB - String message = "a".repeat(messageSize * 1000); - partProducer.newMessage().value(message).send(); - Message msg = consumer.receive(5, TimeUnit.SECONDS); - assertNotNull(msg); - assertTrue(msg.getMessageId() instanceof ChunkMessageIdImpl); - assertEquals(msg.getValue(), message); - - Field msgIdGenerator = ProducerImpl.class.getDeclaredField("msgIdGenerator"); - msgIdGenerator.setAccessible(true); - assertEquals(msg.getSequenceId() + 1, msgIdGenerator.get(partProducer)); - - String message2 = "b".repeat(messageSize * 2); - partProducer.newMessage().value(message2).send(); - Message msg2 = consumer.receive(5, TimeUnit.SECONDS); - assertFalse(msg2.getMessageId() instanceof ChunkMessageIdImpl); - assertEquals(msg2.getValue(), message2); - - long sequenceID = (long) msgIdGenerator.get(partProducer) + 1024L; - String message3 = "c".repeat(messageSize * 1000); - partProducer.newMessage().value(message3).sequenceId(sequenceID).send(); - Message msg3 = consumer.receive(5, TimeUnit.SECONDS); - assertNotNull(msg3); - assertTrue(msg3.getMessageId() instanceof ChunkMessageIdImpl); - assertEquals(msg3.getValue(), message3); - assertEquals(msg3.getSequenceId(), sequenceID); - } - - @Test - public void testDeduplicateChunksInSingleChunkMessages() throws Exception { - this.conf.setBrokerDeduplicationEnabled(true); - restartBroker(); - String topicName = "persistent://my-property/my-ns/testDeduplicateChunksInSingleChunkMessage"; - String producerName = "test-producer"; - @Cleanup - Consumer consumer = pulsarClient - .newConsumer(Schema.STRING) - .subscriptionName("test-sub") - .topic(topicName) - .subscribe(); - final PersistentTopic persistentTopic = (PersistentTopic) pulsar.getBrokerService() - .getTopicIfExists(topicName).get().orElse(null); - assertNotNull(persistentTopic); - sendChunk(persistentTopic, producerName, 1, 0, 2); - sendChunk(persistentTopic, producerName, 1, 1, 2); - sendChunk(persistentTopic, producerName, 1, 1, 2); - - Message message = consumer.receive(15, TimeUnit.SECONDS); - assertEquals(message.getData().length, 2); - - sendChunk(persistentTopic, producerName, 2, 0, 3); - sendChunk(persistentTopic, producerName, 2, 1, 3); - sendChunk(persistentTopic, producerName, 2, 1, 3); - sendChunk(persistentTopic, producerName, 2, 2, 3); - message = consumer.receive(20, TimeUnit.SECONDS); - assertEquals(message.getData().length, 3); - } - private static void sendNonChunk(final PersistentTopic persistentTopic, final String producerName, final long sequenceId) { sendChunk(persistentTopic, producerName, sequenceId, null, null); } - private static void sendChunk(final PersistentTopic persistentTopic, + protected static void sendChunk(final PersistentTopic persistentTopic, final String producerName, final long sequenceId, final Integer chunkId, From 233c5cc8804d99324fd4be8eedfaff2030036967 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Wed, 30 Aug 2023 23:14:52 +0800 Subject: [PATCH 24/24] address some comments --- .../MessageChunkingDeduplicationTest.java | 48 ------------------- .../pulsar/client/impl/ConsumerImpl.java | 7 +-- 2 files changed, 4 insertions(+), 51 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingDeduplicationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingDeduplicationTest.java index 24aa2356017eb..5e590414132a5 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingDeduplicationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingDeduplicationTest.java @@ -20,11 +20,9 @@ import static org.apache.pulsar.client.impl.MessageChunkingSharedTest.sendChunk; import static org.testng.Assert.assertEquals; -import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertNull; import static org.testng.Assert.assertTrue; -import java.lang.reflect.Field; import java.util.concurrent.TimeUnit; import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; @@ -86,52 +84,6 @@ public void testSendChunkMessageWithSameSequenceID() throws Exception { assertNull(msg); } - @Test - public void testDuplicateForChunkMessage() throws Exception { - String topicName = "persistent://my-property/my-ns/testDuplicateForChunkMessage"; - String producerName = "test-producer"; - @Cleanup - Consumer consumer = pulsarClient - .newConsumer(Schema.STRING) - .subscriptionName("test-sub") - .topic(topicName) - .subscribe(); - @Cleanup - Producer partProducer = pulsarClient - .newProducer(Schema.STRING) - .producerName(producerName) - .topic(topicName) - .enableChunking(true) - .enableBatching(false) - .create(); - int messageSize = 6000; // payload size in KB - String message = "a".repeat(messageSize * 1000); - partProducer.newMessage().value(message).send(); - Message msg = consumer.receive(5, TimeUnit.SECONDS); - assertNotNull(msg); - assertTrue(msg.getMessageId() instanceof ChunkMessageIdImpl); - assertEquals(msg.getValue(), message); - - Field msgIdGenerator = ProducerImpl.class.getDeclaredField("msgIdGenerator"); - msgIdGenerator.setAccessible(true); - assertEquals(msg.getSequenceId() + 1, msgIdGenerator.get(partProducer)); - - String message2 = "b".repeat(messageSize * 2); - partProducer.newMessage().value(message2).send(); - Message msg2 = consumer.receive(5, TimeUnit.SECONDS); - assertFalse(msg2.getMessageId() instanceof ChunkMessageIdImpl); - assertEquals(msg2.getValue(), message2); - - long sequenceID = (long) msgIdGenerator.get(partProducer) + 1024L; - String message3 = "c".repeat(messageSize * 1000); - partProducer.newMessage().value(message3).sequenceId(sequenceID).send(); - Message msg3 = consumer.receive(5, TimeUnit.SECONDS); - assertNotNull(msg3); - assertTrue(msg3.getMessageId() instanceof ChunkMessageIdImpl); - assertEquals(msg3.getValue(), message3); - assertEquals(msg3.getSequenceId(), sequenceID); - } - @Test public void testDeduplicateChunksInSingleChunkMessages() throws Exception { String topicName = "persistent://my-property/my-ns/testDeduplicateChunksInSingleChunkMessage"; 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 58b97bf158959..be7c094318434 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 @@ -1470,9 +1470,10 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m // Chunk-6 sequence ID: 0, chunk ID: 3, msgID: 1:6 // We should filter and ack chunk-4 and chunk-5. if (chunkedMsgCtx != null && msgMetadata.getChunkId() <= chunkedMsgCtx.lastChunkedMessageId) { - log.warn("[{}] Receive a repeated chunk messageId {}, last-chunk-id{}, chunkId = {}", - msgMetadata.getProducerName(), chunkedMsgCtx.lastChunkedMessageId, - msgId, msgMetadata.getChunkId()); + log.warn("[{}] Receive a duplicated chunk message with messageId [{}], last-chunk-Id [{}], " + + "chunkId [{}], sequenceId [{}]", + msgMetadata.getProducerName(), msgId, chunkedMsgCtx.lastChunkedMessageId, + msgMetadata.getChunkId(), msgMetadata.getSequenceId()); compressedPayload.release(); increaseAvailablePermits(cnx); boolean repeatedlyReceived = Arrays.stream(chunkedMsgCtx.chunkedMessageIds)