From a4669e07669ea8e9553eddef99a5053130e4f254 Mon Sep 17 00:00:00 2001 From: lipenghui Date: Tue, 29 Oct 2019 17:15:59 +0800 Subject: [PATCH 01/14] Fix message deduplicate issue while using external sequence id with batch --- .../pulsar/broker/service/Producer.java | 61 ++++- .../pulsar/broker/service/ServerCnx.java | 8 +- .../apache/pulsar/broker/service/Topic.java | 8 + .../persistent/MessageDeduplication.java | 28 ++- .../persistent/MessageDuplicationTest.java | 42 ++++ .../client/api/ClientDeduplicationTest.java | 118 ++++++++- .../impl/BatchMessageContainerImpl.java | 20 +- .../pulsar/client/impl/ProducerImpl.java | 56 ++++- .../pulsar/common/api/proto/PulsarApi.java | 228 ++++++++++++++++++ .../pulsar/common/protocol/Commands.java | 33 +++ pulsar-common/src/main/proto/PulsarApi.proto | 8 + 11 files changed, 585 insertions(+), 25 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java index 915c42d3c37dc..d2e0c0271bbe3 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java @@ -132,6 +132,25 @@ public boolean equals(Object obj) { } public void publishMessage(long producerId, long sequenceId, ByteBuf headersAndPayload, long batchSize) { + beforePublish(producerId, sequenceId, headersAndPayload, batchSize); + publishMessageToTopic(headersAndPayload, sequenceId, batchSize); + } + + public void publishMessage(long producerId, long lowestSequenceId, long highestSequenceId, + ByteBuf headersAndPayload, long batchSize) { + if (lowestSequenceId > highestSequenceId) { + cnx.ctx().channel().eventLoop().execute(() -> { + cnx.ctx().writeAndFlush(Commands.newSendError(producerId, highestSequenceId, ServerError.MetadataError, + "Invalid lowest or highest sequence id")); + cnx.completedSendOperation(isNonPersistentTopic); + }); + return; + } + beforePublish(producerId, highestSequenceId, headersAndPayload, batchSize); + publishMessageToTopic(headersAndPayload, lowestSequenceId, highestSequenceId, batchSize); + } + + public void beforePublish(long producerId, long sequenceId, ByteBuf headersAndPayload, long batchSize) { if (isClosed) { cnx.ctx().channel().eventLoop().execute(() -> { cnx.ctx().writeAndFlush(Commands.newSendError(producerId, sequenceId, ServerError.PersistenceError, @@ -170,11 +189,20 @@ public void publishMessage(long producerId, long sequenceId, ByteBuf headersAndP } startPublishOperation((int) batchSize, headersAndPayload.readableBytes()); + } + + private void publishMessageToTopic(ByteBuf headersAndPayload, long sequenceId, long batchSize) { topic.publishMessage(headersAndPayload, MessagePublishContext.get(this, sequenceId, msgIn, headersAndPayload.readableBytes(), batchSize, System.nanoTime())); } + private void publishMessageToTopic(ByteBuf headersAndPayload, long lowestSequenceId, long highestSequenceId, long batchSize) { + topic.publishMessage(headersAndPayload, + MessagePublishContext.get(this, lowestSequenceId, highestSequenceId, msgIn, headersAndPayload.readableBytes(), batchSize, + System.nanoTime())); + } + private boolean verifyChecksum(ByteBuf headersAndPayload) { if (hasChecksum(headersAndPayload)) { int readerIndex = headersAndPayload.readerIndex(); @@ -257,6 +285,9 @@ private static final class MessagePublishContext implements PublishContext, Runn private String originalProducerName; private long originalSequenceId; + private long lowestSequenceId; + private long highestSequenceId; + public String getProducerName() { return producer.getProducerName(); } @@ -265,6 +296,15 @@ public long getSequenceId() { return sequenceId; } + public long getLowestSequenceId() { + return lowestSequenceId; + } + + @Override + public long getHighestSequenceId() { + return highestSequenceId; + } + @Override public void setOriginalProducerName(String originalProducerName) { this.originalProducerName = originalProducerName; @@ -298,7 +338,8 @@ public void completed(Exception exception, long ledgerId, long entryId) { if (!(exception instanceof TopicClosedException)) { // For TopicClosed exception there's no need to send explicit error, since the client was // already notified - producer.cnx.ctx().writeAndFlush(Commands.newSendError(producer.producerId, sequenceId, + long callBackSequenceId = highestSequenceId > 0 ? highestSequenceId : sequenceId; + producer.cnx.ctx().writeAndFlush(Commands.newSendError(producer.producerId, callBackSequenceId, serverError, exception.getMessage())); } producer.cnx.completedSendOperation(producer.isNonPersistentTopic); @@ -330,8 +371,9 @@ public void run() { // stats rateIn.recordMultipleEvents(batchSize, msgSize); producer.topic.recordAddLatency(System.nanoTime() - startTimeNs, TimeUnit.NANOSECONDS); + long callBackSequenceId = highestSequenceId > 0 ? highestSequenceId : sequenceId; producer.cnx.ctx().writeAndFlush( - Commands.newSendReceipt(producer.producerId, sequenceId, ledgerId, entryId), + Commands.newSendReceipt(producer.producerId, callBackSequenceId, ledgerId, entryId), producer.cnx.ctx().voidPromise()); producer.cnx.completedSendOperation(producer.isNonPersistentTopic); producer.publishOperationCompleted(); @@ -352,6 +394,21 @@ static MessagePublishContext get(Producer producer, long sequenceId, Rate rateIn return callback; } + static MessagePublishContext get(Producer producer, long lowestSequenceId, long highestSequenceId, Rate rateIn, + int msgSize, long batchSize, long startTimeNs) { + MessagePublishContext callback = RECYCLER.get(); + callback.producer = producer; + callback.lowestSequenceId = lowestSequenceId; + callback.highestSequenceId = highestSequenceId; + callback.rateIn = rateIn; + callback.msgSize = msgSize; + callback.batchSize = batchSize; + callback.originalProducerName = null; + callback.originalSequenceId = -1; + callback.startTimeNs = startTimeNs; + return callback; + } + private final Handle recyclerHandle; private MessagePublishContext(Handle recyclerHandle) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index d0ba3b93cd907..ec13c8a80b5a3 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -1051,7 +1051,13 @@ protected void handleSend(CommandSend send, ByteBuf headersAndPayload) { startSendOperation(producer); // Persist the message - producer.publishMessage(send.getProducerId(), send.getSequenceId(), headersAndPayload, send.getNumMessages()); + if (send.getNumMessages() > 1 && send.hasLowestSequenceId() && send.getLowestSequenceId() > 0 + && send.hasHighestSequenceId() && send.getHighestSequenceId() > 0) { + producer.publishMessage(send.getProducerId(), send.getLowestSequenceId(), send.getHighestSequenceId(), + headersAndPayload, send.getNumMessages()); + } else { + producer.publishMessage(send.getProducerId(), send.getSequenceId(), headersAndPayload, send.getNumMessages()); + } } private void printSendCommandDebug(CommandSend send, ByteBuf headersAndPayload) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java index 419a9838bcf28..e668562bbd186 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java @@ -77,6 +77,14 @@ default long getOriginalSequenceId() { } void completed(Exception e, long ledgerId, long entryId); + + default long getLowestSequenceId() { + return -1; + } + + default long getHighestSequenceId() { + return -1; + } } void publishMessage(ByteBuf headersAndPayload, PublishContext callback); 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 e898b4ab203cf..ed7814bcf32e1 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 @@ -160,7 +160,8 @@ public void readEntriesComplete(List entries, Object ctx) { MessageMetadata md = Commands.parseMessageMetadata(messageMetadataAndPayload); String producerName = md.getProducerName(); - long sequenceId = md.getSequenceId(); + long sequenceId = md.hasHighestSequenceId() && md.getHighestSequenceId() > 0 ? + md.getHighestSequenceId() : md.getSequenceId(); highestSequencedPushed.put(producerName, sequenceId); highestSequencedPersisted.put(producerName, sequenceId); @@ -283,6 +284,8 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade String producerName = publishContext.getProducerName(); long sequenceId = publishContext.getSequenceId(); + long lowestSequenceId = publishContext.getLowestSequenceId(); + long highestSequenceId = publishContext.getHighestSequenceId(); 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. @@ -290,17 +293,24 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade MessageMetadata md = Commands.parseMessageMetadata(headersAndPayload); producerName = md.getProducerName(); sequenceId = md.getSequenceId(); + lowestSequenceId = md.getLowestSequenceId(); + highestSequenceId = md.getHighestSequenceId(); publishContext.setOriginalProducerName(producerName); publishContext.setOriginalSequenceId(sequenceId); headersAndPayload.readerIndex(readerIndex); md.recycle(); } + if (lowestSequenceId <= highestSequenceId && lowestSequenceId <= 0) { + lowestSequenceId = sequenceId; + highestSequenceId = sequenceId; + } + // 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) { + if (lastSequenceIdPushed != null && lowestSequenceId <= lastSequenceIdPushed) { if (log.isDebugEnabled()) { log.debug("[{}] Message identified as duplicated producer={} seq-id={} -- highest-seq-id={}", topic.getName(), producerName, sequenceId, lastSequenceIdPushed); @@ -311,14 +321,14 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade // If current message's seq id is between lastSequenceIdPersisted and 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 && lowestSequenceId <= lastSequenceIdPersisted) { return MessageDupStatus.Dup; } else { return MessageDupStatus.Unknown; } } - highestSequencedPushed.put(producerName, sequenceId); + highestSequencedPushed.put(producerName, highestSequenceId); } return MessageDupStatus.NotDup; } @@ -333,13 +343,21 @@ public void recordMessagePersisted(PublishContext publishContext, PositionImpl p String producerName = publishContext.getProducerName(); long sequenceId = publishContext.getSequenceId(); + long lowestSequenceId = publishContext.getLowestSequenceId(); + long highestSequenceId = publishContext.getHighestSequenceId(); if (publishContext.getOriginalProducerName() != null) { // In case of replicated messages, this will be different from the current replicator producer name producerName = publishContext.getOriginalProducerName(); sequenceId = publishContext.getOriginalSequenceId(); + lowestSequenceId = publishContext.getLowestSequenceId(); + highestSequenceId = publishContext.getHighestSequenceId(); + } + + if (lowestSequenceId <= highestSequenceId && lowestSequenceId <= 0) { + highestSequenceId = sequenceId; } - highestSequencedPersisted.put(producerName, sequenceId); + highestSequencedPersisted.put(producerName, highestSequenceId); if (++snapshotCounter >= snapshotInterval) { snapshotCounter = 0; takeSnapshot(position); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java index a29de119b0cdd..0b205ddf505e0 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java @@ -45,6 +45,7 @@ import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertTrue; @Slf4j @@ -125,6 +126,24 @@ public void testIsDuplicate() { lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(producerName1); assertTrue(lastSequenceIdPushed != null); assertEquals(lastSequenceIdPushed.longValue(), 5); + + // update highest sequence persisted + messageDeduplication.highestSequencedPushed.put(producerName1, 0L); + messageDeduplication.highestSequencedPersisted.put(producerName1, 0L); + byteBuf1 = getMessage(producerName1, 0); + publishContext1 = getPublishContext(producerName1, 1, 5); + status = messageDeduplication.isDuplicate(publishContext1, byteBuf1); + assertEquals(status, MessageDeduplication.MessageDupStatus.NotDup); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(producerName1); + assertNotNull(lastSequenceIdPushed); + assertEquals(lastSequenceIdPushed.longValue(), 5); + + publishContext1 = getPublishContext(producerName1, 4, 8); + status = messageDeduplication.isDuplicate(publishContext1, byteBuf1); + assertEquals(status, MessageDeduplication.MessageDupStatus.Unknown); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(producerName1); + assertNotNull(lastSequenceIdPushed); + assertEquals(lastSequenceIdPushed.longValue(), 5); } @Test @@ -319,4 +338,27 @@ public void completed(Exception e, long ledgerId, long entryId) { } }); } + + public Topic.PublishContext getPublishContext(String producerName, long lowestSequenceId, long highestSequenceId) { + return spy(new Topic.PublishContext() { + @Override + public String getProducerName() { + return producerName; + } + + public long getLowestSequenceId() { + return lowestSequenceId; + } + + @Override + public long getHighestSequenceId() { + return highestSequenceId; + } + + @Override + public void completed(Exception e, long ledgerId, long entryId) { + + } + }); + } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ClientDeduplicationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ClientDeduplicationTest.java index 03f12b42486f9..916873ebeb0a0 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ClientDeduplicationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ClientDeduplicationTest.java @@ -20,7 +20,10 @@ import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNull; +import static org.testng.Assert.assertTrue; +import static org.testng.Assert.fail; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; import org.testng.annotations.AfterClass; @@ -126,8 +129,12 @@ public void testProducerDeduplication() throws Exception { producer.newMessage().value("my-message-2".getBytes()).sequenceId(2).send(); // Repeat the messages and verify they're not received by consumer - producer.newMessage().value("my-message-1".getBytes()).sequenceId(1).send(); - producer.newMessage().value("my-message-2".getBytes()).sequenceId(2).send(); + try { + producer.newMessage().value("my-message-1".getBytes()).sequenceId(1).send(); + producer.newMessage().value("my-message-2".getBytes()).sequenceId(2).send(); + fail("should be failed"); + } catch (PulsarClientException ignore) { + } producer.close(); @@ -156,4 +163,111 @@ public void testProducerDeduplication() throws Exception { producer.close(); } + + @Test(timeOut = 30000) + public void testProducerDeduplicationWithDiscontinuousSequenceId() throws Exception { + String topic = "persistent://my-property/my-ns/testProducerDeduplicationWithDiscontinuousSequenceId"; + admin.namespaces().setDeduplicationStatus("my-property/my-ns", true); + + // Set infinite timeout + ProducerBuilder producerBuilder = pulsarClient.newProducer().topic(topic) + .producerName("my-producer-name").enableBatching(true).batchingMaxMessages(10).sendTimeout(0, TimeUnit.SECONDS); + Producer producer = producerBuilder.create(); + + assertEquals(producer.getLastSequenceId(), -1L); + + Consumer consumer = pulsarClient.newConsumer().topic(topic).subscriptionName("my-subscription") + .subscribe(); + + producer.newMessage().value("my-message-0".getBytes()).sequenceId(2).sendAsync(); + producer.newMessage().value("my-message-1".getBytes()).sequenceId(3).sendAsync(); + producer.newMessage().value("my-message-2".getBytes()).sequenceId(5).sendAsync(); + + producer.flush(); + + // Repeat the messages and verify they're not received by consumer + CompletableFuture sendResult = producer.newMessage().value("my-message-1".getBytes()).sequenceId(2).sendAsync(); + assertTrue(sendResult.isCompletedExceptionally()); + sendResult = producer.newMessage().value("my-message-2".getBytes()).sequenceId(4).sendAsync(); + assertTrue(sendResult.isCompletedExceptionally()); + producer.close(); + + for (int i = 0; i < 3; i++) { + Message msg = consumer.receive(); + assertEquals(new String(msg.getData()), "my-message-" + i); + consumer.acknowledge(msg); + } + + // No other messages should be received + Message msg = consumer.receive(1, TimeUnit.SECONDS); + assertNull(msg); + + // Kill and restart broker + restartBroker(); + + producer = producerBuilder.create(); + assertEquals(producer.getLastSequenceId(), 5L); + + // Repeat the messages and verify they're not received by consumer + producer.newMessage().value("my-message-1".getBytes()).sequenceId(2).sendAsync(); + producer.newMessage().value("my-message-2".getBytes()).sequenceId(4).sendAsync(); + producer.flush(); + + msg = consumer.receive(1, TimeUnit.SECONDS); + assertNull(msg); + + producer.close(); + } + + @Test(timeOut = 30000) + public void testProducerDeduplicationNonBatchAsync() throws Exception { + String topic = "persistent://my-property/my-ns/testProducerDeduplicationNonBatchAsync"; + admin.namespaces().setDeduplicationStatus("my-property/my-ns", true); + + // Set infinite timeout + ProducerBuilder producerBuilder = pulsarClient.newProducer().topic(topic) + .producerName("my-producer-name").enableBatching(false).sendTimeout(0, TimeUnit.SECONDS); + Producer producer = producerBuilder.create(); + + assertEquals(producer.getLastSequenceId(), -1L); + + Consumer consumer = pulsarClient.newConsumer().topic(topic).subscriptionName("my-subscription") + .subscribe(); + + producer.newMessage().value("my-message-0".getBytes()).sequenceId(2).sendAsync(); + producer.newMessage().value("my-message-1".getBytes()).sequenceId(3).sendAsync(); + producer.newMessage().value("my-message-2".getBytes()).sequenceId(5).sendAsync(); + + // Repeat the messages and verify they're not received by consumer + CompletableFuture sendResult = producer.newMessage().value("my-message-1".getBytes()).sequenceId(2).sendAsync(); + assertTrue(sendResult.isCompletedExceptionally()); + sendResult = producer.newMessage().value("my-message-2".getBytes()).sequenceId(4).sendAsync(); + assertTrue(sendResult.isCompletedExceptionally()); + producer.close(); + + for (int i = 0; i < 3; i++) { + Message msg = consumer.receive(); + assertEquals(new String(msg.getData()), "my-message-" + i); + consumer.acknowledge(msg); + } + + // No other messages should be received + Message msg = consumer.receive(1, TimeUnit.SECONDS); + assertNull(msg); + + // Kill and restart broker + restartBroker(); + + producer = producerBuilder.create(); + assertEquals(producer.getLastSequenceId(), 5L); + + // Repeat the messages and verify they're not received by consumer + producer.newMessage().value("my-message-1".getBytes()).sequenceId(2).sendAsync(); + producer.newMessage().value("my-message-2".getBytes()).sequenceId(4).sendAsync(); + + msg = consumer.receive(1, TimeUnit.SECONDS); + assertNull(msg); + + producer.close(); + } } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchMessageContainerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchMessageContainerImpl.java index b5ee767304bb6..cd093f7041a6d 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchMessageContainerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchMessageContainerImpl.java @@ -50,7 +50,8 @@ class BatchMessageContainerImpl extends AbstractBatchMessageContainer { private PulsarApi.MessageMetadata.Builder messageMetadata = PulsarApi.MessageMetadata.newBuilder(); // sequence id for this batch which will be persisted as a single entry by broker - private long sequenceId = -1; + private long lowestSequenceId = -1; + private long highestSequenceId = -1; private ByteBuf batchedMessageMetadataAndPayload; private List> messages = Lists.newArrayList(); protected SendCallback previousCallback = null; @@ -68,7 +69,7 @@ public void add(MessageImpl msg, SendCallback callback) { if (++numMessagesInBatch == 1) { // some properties are common amongst the different messages in the batch, hence we just pick it up from // the first message - sequenceId = Commands.initBatchMessageMetadata(messageMetadata, msg.getMessageBuilder()); + lowestSequenceId = Commands.initBatchMessageMetadata(messageMetadata, msg.getMessageBuilder()); this.firstCallback = callback; batchedMessageMetadataAndPayload = PulsarByteBufAllocator.DEFAULT .buffer(Math.min(maxBatchSize, MAX_MESSAGE_BATCH_SIZE_BYTES)); @@ -80,6 +81,8 @@ public void add(MessageImpl msg, SendCallback callback) { previousCallback = callback; currentBatchSizeBytes += msg.getDataBuffer().readableBytes(); messages.add(msg); + highestSequenceId = msg.getSequenceId(); + producer.lastSequenceIdPushed = msg.getSequenceId(); } private ByteBuf getCompressedBatchMetadataAndPayload() { @@ -131,7 +134,8 @@ public void clear() { messageMetadata.clear(); numMessagesInBatch = 0; currentBatchSizeBytes = 0; - sequenceId = -1; + lowestSequenceId = -1; + highestSequenceId = -1; batchedMessageMetadataAndPayload = null; } @@ -149,7 +153,7 @@ public void discard(Exception ex) { } } catch (Throwable t) { log.warn("[{}] [{}] Got exception while completing the callback for msg {}:", topicName, producerName, - sequenceId, t); + lowestSequenceId, t); } ReferenceCountUtil.safeRelease(batchedMessageMetadataAndPayload); clear(); @@ -164,10 +168,12 @@ public boolean isMultiBatches() { public OpSendMsg createOpSendMsg() throws IOException { ByteBuf encryptedPayload = producer.encryptMessage(messageMetadata, getCompressedBatchMetadataAndPayload()); messageMetadata.setNumMessagesInBatch(numMessagesInBatch); - ByteBufPair cmd = producer.sendMessage(producer.producerId, sequenceId, numMessagesInBatch, - messageMetadata.build(), encryptedPayload); + messageMetadata.setHighestSequenceId(highestSequenceId); + ByteBufPair cmd = producer.sendMessage(producer.producerId, messageMetadata.getLowestSequenceId(), + messageMetadata.getHighestSequenceId(), numMessagesInBatch, messageMetadata.build(), encryptedPayload); - OpSendMsg op = OpSendMsg.create(messages, cmd, sequenceId, firstCallback); + OpSendMsg op = OpSendMsg.create(messages, cmd, messageMetadata.getLowestSequenceId(), + messageMetadata.getHighestSequenceId(), firstCallback); if (encryptedPayload.readableBytes() > ClientCnx.getMaxMessageSize()) { cmd.release(); 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 ce14271ee0c30..6be2f390b8084 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 @@ -110,6 +110,8 @@ public class ProducerImpl extends ProducerBase implements TimerTask, Conne private final CompressionCodec compressor; private volatile long lastSequenceIdPublished; + protected volatile long lastSequenceIdPushed; + private MessageCrypto msgCrypto = null; private ScheduledFuture keyGeneratorTask = null; @@ -140,9 +142,11 @@ public ProducerImpl(PulsarClientImpl client, String topic, ProducerConfiguration if (conf.getInitialSequenceId() != null) { long initialSequenceId = conf.getInitialSequenceId(); this.lastSequenceIdPublished = initialSequenceId; + this.lastSequenceIdPushed = initialSequenceId; this.msgIdGenerator = initialSequenceId + 1; } else { this.lastSequenceIdPublished = -1; + this.lastSequenceIdPushed = -1; this.msgIdGenerator = 0; } @@ -368,6 +372,19 @@ public void sendAsync(Message message, SendCallback callback) { } else { sequenceId = msgMetadataBuilder.getSequenceId(); } + + if (sequenceId <= lastSequenceIdPushed) { + if (sequenceId <= lastSequenceIdPublished) { + callback.sendComplete(new PulsarClientException + .InvalidMessageException("Message is definitely a duplicate")); + return; + } else { + callback.sendComplete(new PulsarClientException + .InvalidMessageException("Message is a definitely a duplicate or not cannot be " + + "determined at this time")); + return; + } + } if (!msgMetadataBuilder.hasPublishTime()) { msgMetadataBuilder.setPublishTime(client.getClientClock().millis()); @@ -534,15 +551,21 @@ protected ByteBuf encryptMessage(MessageMetadata.Builder msgMetadata, ByteBuf co protected ByteBufPair sendMessage(long producerId, long sequenceId, int numMessages, MessageMetadata msgMetadata, ByteBuf compressedPayload) { - ChecksumType checksumType; + return Commands.newSend(producerId, sequenceId, numMessages, getChecksumType(), msgMetadata, compressedPayload); + } + protected ByteBufPair sendMessage(long producerId, long lowestSequenceId, long highestSequenceId, int numMessages, MessageMetadata msgMetadata, + ByteBuf compressedPayload) { + return Commands.newSend(producerId, lowestSequenceId, highestSequenceId, numMessages, getChecksumType(), msgMetadata, compressedPayload); + } + + private ChecksumType getChecksumType() { if (connectionHandler.getClientCnx() == null || connectionHandler.getClientCnx().getRemoteEndpointProtocolVersion() >= brokerChecksumSupportedVersion()) { - checksumType = ChecksumType.Crc32c; + return ChecksumType.Crc32c; } else { - checksumType = ChecksumType.None; + return ChecksumType.None; } - return Commands.newSend(producerId, sequenceId, numMessages, checksumType, msgMetadata, compressedPayload); } private boolean canAddToBatch(MessageImpl msg) { @@ -776,8 +799,7 @@ void ackReceived(ClientCnx cnx, long sequenceId, long ledgerId, long entryId) { } return; } - - long expectedSequenceId = op.sequenceId; + long expectedSequenceId = op.highestSequenceId > 0 ? op.highestSequenceId : op.sequenceId; if (sequenceId > expectedSequenceId) { log.warn("[{}] [{}] Got ack for msg. expecting: {} - got: {} - queue-size: {}", topic, producerName, expectedSequenceId, sequenceId, pendingMessages.size()); @@ -803,7 +825,7 @@ void ackReceived(ClientCnx cnx, long sequenceId, long ledgerId, long entryId) { if (callback) { op = pendingCallbacks.poll(); if (op != null) { - lastSequenceIdPublished = op.sequenceId + op.numMessagesInBatch - 1; + lastSequenceIdPublished = sequenceId; op.setMessageId(ledgerId, entryId, partitionIndex); try { // Need to protect ourselves from any exception being thrown in the future handler from the @@ -839,7 +861,7 @@ protected synchronized void recoverChecksumError(ClientCnx cnx, long sequenceId) log.debug("[{}] [{}] Got send failure for timed out msg {}", topic, producerName, sequenceId); } } else { - long expectedSequenceId = op.sequenceId; + long expectedSequenceId = op.highestSequenceId > 0 ? op.highestSequenceId : op.sequenceId; if (sequenceId == expectedSequenceId) { boolean corrupted = !verifyLocalBufferIsNotCorrupted(op); if (corrupted) { @@ -923,6 +945,8 @@ protected static final class OpSendMsg { long createdAt; long batchSizeByte = 0; int numMessagesInBatch = 1; + long lowestSequenceId; + long highestSequenceId; static OpSendMsg create(MessageImpl msg, ByteBufPair cmd, long sequenceId, SendCallback callback) { OpSendMsg op = RECYCLER.get(); @@ -944,6 +968,19 @@ static OpSendMsg create(List> msgs, ByteBufPair cmd, long sequenc return op; } + static OpSendMsg create(List> msgs, ByteBufPair cmd, long lowestSequenceId, + long highestSequenceId, SendCallback callback) { + OpSendMsg op = RECYCLER.get(); + op.msgs = msgs; + op.cmd = cmd; + op.callback = callback; + op.lowestSequenceId = lowestSequenceId; + op.highestSequenceId = highestSequenceId; + op.sequenceId = lowestSequenceId; + op.createdAt = System.currentTimeMillis(); + return op; + } + void recycle() { msg = null; msgs = null; @@ -1391,6 +1428,9 @@ private void processOpSendMsg(OpSendMsg op) { batchMessageAndSend(); } pendingMessages.put(op); + if (op.msg != null) { + lastSequenceIdPushed = op.sequenceId; + } ClientCnx cnx = cnx(); if (isConnected()) { if (op.msg != null && op.msg.getSchemaState() == None) { diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java b/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java index 2074659432d05..145051ace9654 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java @@ -3589,6 +3589,14 @@ public interface MessageMetadataOrBuilder // optional uint64 txnid_most_bits = 23 [default = 0]; boolean hasTxnidMostBits(); long getTxnidMostBits(); + + // optional uint64 lowest_sequence_id = 24 [default = 0]; + boolean hasLowestSequenceId(); + long getLowestSequenceId(); + + // optional uint64 highest_sequence_id = 25 [default = 0]; + boolean hasHighestSequenceId(); + long getHighestSequenceId(); } public static final class MessageMetadata extends org.apache.pulsar.shaded.com.google.protobuf.v241.GeneratedMessageLite @@ -3949,6 +3957,26 @@ public long getTxnidMostBits() { return txnidMostBits_; } + // optional uint64 lowest_sequence_id = 24 [default = 0]; + public static final int LOWEST_SEQUENCE_ID_FIELD_NUMBER = 24; + private long lowestSequenceId_; + public boolean hasLowestSequenceId() { + return ((bitField0_ & 0x00040000) == 0x00040000); + } + public long getLowestSequenceId() { + return lowestSequenceId_; + } + + // optional uint64 highest_sequence_id = 25 [default = 0]; + public static final int HIGHEST_SEQUENCE_ID_FIELD_NUMBER = 25; + private long highestSequenceId_; + public boolean hasHighestSequenceId() { + return ((bitField0_ & 0x00080000) == 0x00080000); + } + public long getHighestSequenceId() { + return highestSequenceId_; + } + private void initFields() { producerName_ = ""; sequenceId_ = 0L; @@ -3971,6 +3999,8 @@ private void initFields() { markerType_ = 0; txnidLeastBits_ = 0L; txnidMostBits_ = 0L; + lowestSequenceId_ = 0L; + highestSequenceId_ = 0L; } private byte memoizedIsInitialized = -1; public final boolean isInitialized() { @@ -4076,6 +4106,12 @@ public void writeTo(org.apache.pulsar.common.util.protobuf.ByteBufCodedOutputStr if (((bitField0_ & 0x00020000) == 0x00020000)) { output.writeUInt64(23, txnidMostBits_); } + if (((bitField0_ & 0x00040000) == 0x00040000)) { + output.writeUInt64(24, lowestSequenceId_); + } + if (((bitField0_ & 0x00080000) == 0x00080000)) { + output.writeUInt64(25, highestSequenceId_); + } } private int memoizedSerializedSize = -1; @@ -4173,6 +4209,14 @@ public int getSerializedSize() { size += org.apache.pulsar.shaded.com.google.protobuf.v241.CodedOutputStream .computeUInt64Size(23, txnidMostBits_); } + if (((bitField0_ & 0x00040000) == 0x00040000)) { + size += org.apache.pulsar.shaded.com.google.protobuf.v241.CodedOutputStream + .computeUInt64Size(24, lowestSequenceId_); + } + if (((bitField0_ & 0x00080000) == 0x00080000)) { + size += org.apache.pulsar.shaded.com.google.protobuf.v241.CodedOutputStream + .computeUInt64Size(25, highestSequenceId_); + } memoizedSerializedSize = size; return size; } @@ -4328,6 +4372,10 @@ public Builder clear() { bitField0_ = (bitField0_ & ~0x00080000); txnidMostBits_ = 0L; bitField0_ = (bitField0_ & ~0x00100000); + lowestSequenceId_ = 0L; + bitField0_ = (bitField0_ & ~0x00200000); + highestSequenceId_ = 0L; + bitField0_ = (bitField0_ & ~0x00400000); return this; } @@ -4449,6 +4497,14 @@ public org.apache.pulsar.common.api.proto.PulsarApi.MessageMetadata buildPartial to_bitField0_ |= 0x00020000; } result.txnidMostBits_ = txnidMostBits_; + if (((from_bitField0_ & 0x00200000) == 0x00200000)) { + to_bitField0_ |= 0x00040000; + } + result.lowestSequenceId_ = lowestSequenceId_; + if (((from_bitField0_ & 0x00400000) == 0x00400000)) { + to_bitField0_ |= 0x00080000; + } + result.highestSequenceId_ = highestSequenceId_; result.bitField0_ = to_bitField0_; return result; } @@ -4539,6 +4595,12 @@ public Builder mergeFrom(org.apache.pulsar.common.api.proto.PulsarApi.MessageMet if (other.hasTxnidMostBits()) { setTxnidMostBits(other.getTxnidMostBits()); } + if (other.hasLowestSequenceId()) { + setLowestSequenceId(other.getLowestSequenceId()); + } + if (other.hasHighestSequenceId()) { + setHighestSequenceId(other.getHighestSequenceId()); + } return this; } @@ -4703,6 +4765,16 @@ public Builder mergeFrom( txnidMostBits_ = input.readUInt64(); break; } + case 192: { + bitField0_ |= 0x00200000; + lowestSequenceId_ = input.readUInt64(); + break; + } + case 200: { + bitField0_ |= 0x00400000; + highestSequenceId_ = input.readUInt64(); + break; + } } } } @@ -5393,6 +5465,48 @@ public Builder clearTxnidMostBits() { return this; } + // optional uint64 lowest_sequence_id = 24 [default = 0]; + private long lowestSequenceId_ ; + public boolean hasLowestSequenceId() { + return ((bitField0_ & 0x00200000) == 0x00200000); + } + public long getLowestSequenceId() { + return lowestSequenceId_; + } + public Builder setLowestSequenceId(long value) { + bitField0_ |= 0x00200000; + lowestSequenceId_ = value; + + return this; + } + public Builder clearLowestSequenceId() { + bitField0_ = (bitField0_ & ~0x00200000); + lowestSequenceId_ = 0L; + + return this; + } + + // optional uint64 highest_sequence_id = 25 [default = 0]; + private long highestSequenceId_ ; + public boolean hasHighestSequenceId() { + return ((bitField0_ & 0x00400000) == 0x00400000); + } + public long getHighestSequenceId() { + return highestSequenceId_; + } + public Builder setHighestSequenceId(long value) { + bitField0_ |= 0x00400000; + highestSequenceId_ = value; + + return this; + } + public Builder clearHighestSequenceId() { + bitField0_ = (bitField0_ & ~0x00400000); + highestSequenceId_ = 0L; + + return this; + } + // @@protoc_insertion_point(builder_scope:pulsar.proto.MessageMetadata) } @@ -15327,6 +15441,14 @@ public interface CommandSendOrBuilder // optional uint64 txnid_most_bits = 5 [default = 0]; boolean hasTxnidMostBits(); long getTxnidMostBits(); + + // optional uint64 lowest_sequence_id = 6 [default = 0]; + boolean hasLowestSequenceId(); + long getLowestSequenceId(); + + // optional uint64 highest_sequence_id = 7 [default = 0]; + boolean hasHighestSequenceId(); + long getHighestSequenceId(); } public static final class CommandSend extends org.apache.pulsar.shaded.com.google.protobuf.v241.GeneratedMessageLite @@ -15413,12 +15535,34 @@ public long getTxnidMostBits() { return txnidMostBits_; } + // optional uint64 lowest_sequence_id = 6 [default = 0]; + public static final int LOWEST_SEQUENCE_ID_FIELD_NUMBER = 6; + private long lowestSequenceId_; + public boolean hasLowestSequenceId() { + return ((bitField0_ & 0x00000020) == 0x00000020); + } + public long getLowestSequenceId() { + return lowestSequenceId_; + } + + // optional uint64 highest_sequence_id = 7 [default = 0]; + public static final int HIGHEST_SEQUENCE_ID_FIELD_NUMBER = 7; + private long highestSequenceId_; + public boolean hasHighestSequenceId() { + return ((bitField0_ & 0x00000040) == 0x00000040); + } + public long getHighestSequenceId() { + return highestSequenceId_; + } + private void initFields() { producerId_ = 0L; sequenceId_ = 0L; numMessages_ = 1; txnidLeastBits_ = 0L; txnidMostBits_ = 0L; + lowestSequenceId_ = 0L; + highestSequenceId_ = 0L; } private byte memoizedIsInitialized = -1; public final boolean isInitialized() { @@ -15460,6 +15604,12 @@ public void writeTo(org.apache.pulsar.common.util.protobuf.ByteBufCodedOutputStr if (((bitField0_ & 0x00000010) == 0x00000010)) { output.writeUInt64(5, txnidMostBits_); } + if (((bitField0_ & 0x00000020) == 0x00000020)) { + output.writeUInt64(6, lowestSequenceId_); + } + if (((bitField0_ & 0x00000040) == 0x00000040)) { + output.writeUInt64(7, highestSequenceId_); + } } private int memoizedSerializedSize = -1; @@ -15488,6 +15638,14 @@ public int getSerializedSize() { size += org.apache.pulsar.shaded.com.google.protobuf.v241.CodedOutputStream .computeUInt64Size(5, txnidMostBits_); } + if (((bitField0_ & 0x00000020) == 0x00000020)) { + size += org.apache.pulsar.shaded.com.google.protobuf.v241.CodedOutputStream + .computeUInt64Size(6, lowestSequenceId_); + } + if (((bitField0_ & 0x00000040) == 0x00000040)) { + size += org.apache.pulsar.shaded.com.google.protobuf.v241.CodedOutputStream + .computeUInt64Size(7, highestSequenceId_); + } memoizedSerializedSize = size; return size; } @@ -15611,6 +15769,10 @@ public Builder clear() { bitField0_ = (bitField0_ & ~0x00000008); txnidMostBits_ = 0L; bitField0_ = (bitField0_ & ~0x00000010); + lowestSequenceId_ = 0L; + bitField0_ = (bitField0_ & ~0x00000020); + highestSequenceId_ = 0L; + bitField0_ = (bitField0_ & ~0x00000040); return this; } @@ -15664,6 +15826,14 @@ public org.apache.pulsar.common.api.proto.PulsarApi.CommandSend buildPartial() { to_bitField0_ |= 0x00000010; } result.txnidMostBits_ = txnidMostBits_; + if (((from_bitField0_ & 0x00000020) == 0x00000020)) { + to_bitField0_ |= 0x00000020; + } + result.lowestSequenceId_ = lowestSequenceId_; + if (((from_bitField0_ & 0x00000040) == 0x00000040)) { + to_bitField0_ |= 0x00000040; + } + result.highestSequenceId_ = highestSequenceId_; result.bitField0_ = to_bitField0_; return result; } @@ -15685,6 +15855,12 @@ public Builder mergeFrom(org.apache.pulsar.common.api.proto.PulsarApi.CommandSen if (other.hasTxnidMostBits()) { setTxnidMostBits(other.getTxnidMostBits()); } + if (other.hasLowestSequenceId()) { + setLowestSequenceId(other.getLowestSequenceId()); + } + if (other.hasHighestSequenceId()) { + setHighestSequenceId(other.getHighestSequenceId()); + } return this; } @@ -15747,6 +15923,16 @@ public Builder mergeFrom( txnidMostBits_ = input.readUInt64(); break; } + case 48: { + bitField0_ |= 0x00000020; + lowestSequenceId_ = input.readUInt64(); + break; + } + case 56: { + bitField0_ |= 0x00000040; + highestSequenceId_ = input.readUInt64(); + break; + } } } } @@ -15858,6 +16044,48 @@ public Builder clearTxnidMostBits() { return this; } + // optional uint64 lowest_sequence_id = 6 [default = 0]; + private long lowestSequenceId_ ; + public boolean hasLowestSequenceId() { + return ((bitField0_ & 0x00000020) == 0x00000020); + } + public long getLowestSequenceId() { + return lowestSequenceId_; + } + public Builder setLowestSequenceId(long value) { + bitField0_ |= 0x00000020; + lowestSequenceId_ = value; + + return this; + } + public Builder clearLowestSequenceId() { + bitField0_ = (bitField0_ & ~0x00000020); + lowestSequenceId_ = 0L; + + return this; + } + + // optional uint64 highest_sequence_id = 7 [default = 0]; + private long highestSequenceId_ ; + public boolean hasHighestSequenceId() { + return ((bitField0_ & 0x00000040) == 0x00000040); + } + public long getHighestSequenceId() { + return highestSequenceId_; + } + public Builder setHighestSequenceId(long value) { + bitField0_ |= 0x00000040; + highestSequenceId_ = value; + + return this; + } + public Builder clearHighestSequenceId() { + bitField0_ = (bitField0_ & ~0x00000040); + highestSequenceId_ = 0L; + + return this; + } + // @@protoc_insertion_point(builder_scope:pulsar.proto.CommandSend) } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java index bc9ed1ca5cdd6..a1a0c829aa563 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java @@ -443,6 +443,12 @@ public static ByteBufPair newSend(long producerId, long sequenceId, int numMessa return newSend(producerId, sequenceId, numMessaegs, 0, 0, checksumType, messageMetadata, payload); } + public static ByteBufPair newSend(long producerId, long lowestSequenceId, long highestSequenceId, int numMessaegs, ChecksumType checksumType, + MessageMetadata messageMetadata, ByteBuf payload) { + return newSend(producerId, lowestSequenceId, highestSequenceId, numMessaegs, 0, 0, + checksumType, messageMetadata, payload); + } + public static ByteBufPair newSend(long producerId, long sequenceId, int numMessages, long txnIdLeastBits, long txnIdMostBits, ChecksumType checksumType, MessageMetadata messageData, ByteBuf payload) { @@ -467,6 +473,32 @@ public static ByteBufPair newSend(long producerId, long sequenceId, int numMessa return res; } + public static ByteBufPair newSend(long producerId, long lowestSequenceId, long highestSequenceId, int numMessages, + long txnIdLeastBits, long txnIdMostBits, ChecksumType checksumType, + MessageMetadata messageData, ByteBuf payload) { + CommandSend.Builder sendBuilder = CommandSend.newBuilder(); + sendBuilder.setProducerId(producerId); + sendBuilder.setLowestSequenceId(lowestSequenceId); + sendBuilder.setHighestSequenceId(highestSequenceId); + sendBuilder.setSequenceId(lowestSequenceId); + if (numMessages > 1) { + sendBuilder.setNumMessages(numMessages); + } + if (txnIdLeastBits > 0) { + sendBuilder.setTxnidLeastBits(txnIdLeastBits); + } + if (txnIdMostBits > 0) { + sendBuilder.setTxnidMostBits(txnIdMostBits); + } + CommandSend send = sendBuilder.build(); + + ByteBufPair res = serializeCommandSendWithSize(BaseCommand.newBuilder().setType(Type.SEND).setSend(send), + checksumType, messageData, payload); + send.recycle(); + sendBuilder.recycle(); + return res; + } + public static ByteBuf newSubscribe(String topic, String subscription, long consumerId, long requestId, SubType subType, int priorityLevel, String consumerName, long resetStartMessageBackInSeconds) { return newSubscribe(topic, subscription, consumerId, requestId, subType, priorityLevel, consumerName, @@ -1515,6 +1547,7 @@ public static long initBatchMessageMetadata(MessageMetadata.Builder messageMetad messageMetadata.setPublishTime(builder.getPublishTime()); messageMetadata.setProducerName(builder.getProducerName()); messageMetadata.setSequenceId(builder.getSequenceId()); + messageMetadata.setLowestSequenceId(builder.getSequenceId()); if (builder.hasReplicatedFrom()) { messageMetadata.setReplicatedFrom(builder.getReplicatedFrom()); } diff --git a/pulsar-common/src/main/proto/PulsarApi.proto b/pulsar-common/src/main/proto/PulsarApi.proto index 66cb7b58661ba..d76d2c4a33048 100644 --- a/pulsar-common/src/main/proto/PulsarApi.proto +++ b/pulsar-common/src/main/proto/PulsarApi.proto @@ -132,6 +132,10 @@ message MessageMetadata { // transaction related message info optional uint64 txnid_least_bits = 22 [default = 0]; optional uint64 txnid_most_bits = 23 [default = 0]; + + /// Add lowest and highest sequence id to support external sequence id + optional uint64 lowest_sequence_id = 24 [default = 0]; + optional uint64 highest_sequence_id = 25 [default = 0]; } message SingleMessageMetadata { @@ -411,6 +415,10 @@ message CommandSend { optional int32 num_messages = 3 [default = 1]; optional uint64 txnid_least_bits = 4 [default = 0]; optional uint64 txnid_most_bits = 5 [default = 0]; + + /// Add lowest and highest sequence id to support external sequence id + optional uint64 lowest_sequence_id = 6 [default = 0]; + optional uint64 highest_sequence_id = 7 [default = 0]; } message CommandSendReceipt { From 4783321069d1ba60c43eeda8189c3de6863aeda7 Mon Sep 17 00:00:00 2001 From: lipenghui Date: Tue, 29 Oct 2019 17:35:30 +0800 Subject: [PATCH 02/14] Fix checkstyle --- .../main/java/org/apache/pulsar/common/protocol/Commands.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java index a1a0c829aa563..00c35307e5f9d 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java @@ -443,8 +443,8 @@ public static ByteBufPair newSend(long producerId, long sequenceId, int numMessa return newSend(producerId, sequenceId, numMessaegs, 0, 0, checksumType, messageMetadata, payload); } - public static ByteBufPair newSend(long producerId, long lowestSequenceId, long highestSequenceId, int numMessaegs, ChecksumType checksumType, - MessageMetadata messageMetadata, ByteBuf payload) { + public static ByteBufPair newSend(long producerId, long lowestSequenceId, long highestSequenceId, int numMessaegs, + ChecksumType checksumType, MessageMetadata messageMetadata, ByteBuf payload) { return newSend(producerId, lowestSequenceId, highestSequenceId, numMessaegs, 0, 0, checksumType, messageMetadata, payload); } From 5358921c4dfe56a0972bb27000983194e33700b1 Mon Sep 17 00:00:00 2001 From: lipenghui Date: Tue, 29 Oct 2019 21:23:26 +0800 Subject: [PATCH 03/14] Fix unit tests issue --- .../java/org/apache/pulsar/broker/service/ServerCnx.java | 2 +- .../org/apache/pulsar/broker/service/BatchMessageTest.java | 2 ++ .../apache/pulsar/client/impl/BatchMessageContainerImpl.java | 5 +++++ 3 files changed, 8 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index ec13c8a80b5a3..7a4ff34ecaebb 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -1051,7 +1051,7 @@ protected void handleSend(CommandSend send, ByteBuf headersAndPayload) { startSendOperation(producer); // Persist the message - if (send.getNumMessages() > 1 && send.hasLowestSequenceId() && send.getLowestSequenceId() > 0 + if (send.getNumMessages() > 1 && send.hasLowestSequenceId() && send.getLowestSequenceId() >= 0 && send.hasHighestSequenceId() && send.getHighestSequenceId() > 0) { producer.publishMessage(send.getProducerId(), send.getLowestSequenceId(), send.getHighestSequenceId(), headersAndPayload, send.getNumMessages()); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageTest.java index b0ca4b4368d34..f3e024909a769 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageTest.java @@ -567,12 +567,14 @@ public void testBatchAndNonBatchCumulativeAcks(BatcherBuilder builder) throws Ex .batchingMaxMessages(numMsgsInBatch) .enableBatching(true) .batcherBuilder(builder) + .producerName("1") .messageRoutingMode(MessageRoutingMode.SinglePartition) .create(); // create producer to publish non batch messages Producer noBatchProducer = pulsarClient.newProducer().topic(topicName) .enableBatching(false) .messageRoutingMode(MessageRoutingMode.SinglePartition) + .producerName("2") .create(); List> sendFutureList = Lists.newArrayList(); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchMessageContainerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchMessageContainerImpl.java index cd093f7041a6d..ad911d228c80e 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchMessageContainerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchMessageContainerImpl.java @@ -81,6 +81,10 @@ public void add(MessageImpl msg, SendCallback callback) { previousCallback = callback; currentBatchSizeBytes += msg.getDataBuffer().readableBytes(); messages.add(msg); + if (lowestSequenceId == -1) { + lowestSequenceId = msg.getSequenceId(); + messageMetadata.setLowestSequenceId(lowestSequenceId); + } highestSequenceId = msg.getSequenceId(); producer.lastSequenceIdPushed = msg.getSequenceId(); } @@ -187,6 +191,7 @@ public OpSendMsg createOpSendMsg() throws IOException { op.setNumMessagesInBatch(numMessagesInBatch); op.setBatchSizeByte(currentBatchSizeBytes); + lowestSequenceId = -1; return op; } From c8c9d3345726145d3780722e595f7042e579fcfe Mon Sep 17 00:00:00 2001 From: lipenghui Date: Wed, 30 Oct 2019 15:35:21 +0800 Subject: [PATCH 04/14] Fix unit tests issue --- .../main/java/org/apache/pulsar/broker/service/Producer.java | 4 +++- .../main/java/org/apache/pulsar/broker/service/ServerCnx.java | 3 +-- .../java/org/apache/pulsar/client/api/SimpleSchemaTest.java | 2 +- .../main/java/org/apache/pulsar/client/impl/ProducerImpl.java | 2 ++ 4 files changed, 7 insertions(+), 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java index d2e0c0271bbe3..3ee6664b384f0 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java @@ -371,7 +371,7 @@ public void run() { // stats rateIn.recordMultipleEvents(batchSize, msgSize); producer.topic.recordAddLatency(System.nanoTime() - startTimeNs, TimeUnit.NANOSECONDS); - long callBackSequenceId = highestSequenceId > 0 ? highestSequenceId : sequenceId; + long callBackSequenceId = highestSequenceId >= 0 && highestSequenceId >= sequenceId ? highestSequenceId : sequenceId; producer.cnx.ctx().writeAndFlush( Commands.newSendReceipt(producer.producerId, callBackSequenceId, ledgerId, entryId), producer.cnx.ctx().voidPromise()); @@ -424,6 +424,8 @@ protected MessagePublishContext newObject(Recycler.Handle public void recycle() { producer = null; sequenceId = -1; + lowestSequenceId = -1; + highestSequenceId = -1; rateIn = null; msgSize = 0; ledgerId = -1; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index 7a4ff34ecaebb..9b11a92bd3843 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -1051,8 +1051,7 @@ protected void handleSend(CommandSend send, ByteBuf headersAndPayload) { startSendOperation(producer); // Persist the message - if (send.getNumMessages() > 1 && send.hasLowestSequenceId() && send.getLowestSequenceId() >= 0 - && send.hasHighestSequenceId() && send.getHighestSequenceId() > 0) { + if (send.hasLowestSequenceId() && send.hasHighestSequenceId() && send.getSequenceId() <= send.getLowestSequenceId()) { producer.publishMessage(send.getProducerId(), send.getLowestSequenceId(), send.getHighestSequenceId(), headersAndPayload, send.getNumMessages()); } else { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleSchemaTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleSchemaTest.java index aa0f0e0c9018b..5007b2f098444 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleSchemaTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleSchemaTest.java @@ -301,7 +301,7 @@ public void newProducerForMessageOnTopicWithDifferentSchemaType() throws Excepti @Test public void newProducerForMessageSchemaOnTopicInitialWithNoSchema() throws Exception { - String topic = "my-property/my-ns/schema-test"; + String topic = "my-property/my-ns/schema-test-"; Schema v1Schema = Schema.AVRO(V1Data.class); byte[] v1SchemaBytes = v1Schema.getSchemaInfo().getSchema(); AvroWriter v1Writer = new AvroWriter<>( 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 6be2f390b8084..00d9ed5d2a1c1 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 @@ -989,6 +989,8 @@ void recycle() { rePopulate = null; sequenceId = -1; createdAt = -1; + lowestSequenceId = -1; + highestSequenceId = -1; recyclerHandle.recycle(this); } From 4ef11b791a7468b7b8dacf59851f3198951e36d9 Mon Sep 17 00:00:00 2001 From: lipenghui Date: Wed, 30 Oct 2019 17:08:20 +0800 Subject: [PATCH 05/14] Fix unit tests issue --- .../pulsar/client/api/ClientDeduplicationFailureTest.java | 6 +++++- .../test/java/org/apache/pulsar/client/impl/ReaderTest.java | 2 +- 2 files changed, 6 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ClientDeduplicationFailureTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ClientDeduplicationFailureTest.java index 63f8aaa2bcd45..567e9db87973f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ClientDeduplicationFailureTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ClientDeduplicationFailureTest.java @@ -241,7 +241,11 @@ public void testClientDeduplicationCorrectnessWithFailure() throws Exception { producerThread.stop(); // send last message - producer.newMessage().sequenceId(producerThread.getLastSeqId() + 1).value("end").send(); + try { + producer.newMessage().sequenceId(producerThread.getLastSeqId() + 1).value("end").send(); + fail("should failed, because send a duplication"); + } catch (PulsarClientException.InvalidMessageException ignore) { + } producer.close(); Reader reader = pulsarClient.newReader(Schema.STRING).startMessageId(MessageId.earliest).topic(sourceTopic).create(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/ReaderTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/ReaderTest.java index e931c6c2fa017..2919860855a8d 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/ReaderTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/ReaderTest.java @@ -189,7 +189,7 @@ public void testReaderWithTimeLong() throws Exception { TypedMessageBuilderImpl msg = (TypedMessageBuilderImpl) producer.newMessage() .value(("new" + i).getBytes()); Builder metadataBuilder = msg.getMetadataBuilder(); - metadataBuilder.setPublishTime(newMsgPublishTime); + metadataBuilder.setPublishTime(newMsgPublishTime).setSequenceId(totalMsg + i); metadataBuilder.setProducerName(producer.getProducerName()).setReplicatedFrom("us-west1"); MessageId msgId = msg.send(); if (firstMsgId == null) { From b902b376b9c1d4dada445c0f71647bc83b692504 Mon Sep 17 00:00:00 2001 From: lipenghui Date: Wed, 30 Oct 2019 21:48:14 +0800 Subject: [PATCH 06/14] Handle invalid message exception in replicator --- .../pulsar/broker/service/persistent/MessageDeduplication.java | 2 +- .../pulsar/broker/service/persistent/PersistentReplicator.java | 3 ++- 2 files 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 ed7814bcf32e1..88143df38bd77 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 @@ -301,7 +301,7 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade md.recycle(); } - if (lowestSequenceId <= highestSequenceId && lowestSequenceId <= 0) { + if (lowestSequenceId == highestSequenceId && lowestSequenceId <= sequenceId) { lowestSequenceId = sequenceId; highestSequenceId = sequenceId; } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java index 9cc794235dc95..622098e76adec 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java @@ -50,6 +50,7 @@ import org.apache.pulsar.broker.service.Replicator; import org.apache.pulsar.broker.service.persistent.DispatchRateLimiter.Type; import org.apache.pulsar.client.api.MessageId; +import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.impl.Backoff; import org.apache.pulsar.client.impl.MessageImpl; import org.apache.pulsar.client.impl.ProducerImpl; @@ -377,7 +378,7 @@ private static final class ProducerSendCallback implements SendCallback { @Override public void sendComplete(Exception exception) { - if (exception != null) { + if (exception != null && !(exception instanceof PulsarClientException.InvalidMessageException)) { log.error("[{}][{} -> {}] Error producing on remote broker", replicator.topicName, replicator.localCluster, replicator.remoteCluster, exception); // cursor shoud be rewinded since it was incremented when readMoreEntries From f1c1eb2c37c0b766a5d64556199da8dddbb8b168 Mon Sep 17 00:00:00 2001 From: lipenghui Date: Thu, 31 Oct 2019 11:33:26 +0800 Subject: [PATCH 07/14] Use only one last sequence id. --- .../pulsar/broker/service/ServerCnx.java | 4 +- .../persistent/MessageDeduplication.java | 8 +- .../impl/BatchMessageContainerImpl.java | 12 +- .../pulsar/client/impl/ProducerImpl.java | 10 +- .../pulsar/common/api/proto/PulsarApi.java | 218 +++++------------- .../pulsar/common/protocol/Commands.java | 5 +- pulsar-common/src/main/proto/PulsarApi.proto | 10 +- 7 files changed, 71 insertions(+), 196 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index 9b11a92bd3843..5ceb76a5b753e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -1051,8 +1051,8 @@ protected void handleSend(CommandSend send, ByteBuf headersAndPayload) { startSendOperation(producer); // Persist the message - if (send.hasLowestSequenceId() && send.hasHighestSequenceId() && send.getSequenceId() <= send.getLowestSequenceId()) { - producer.publishMessage(send.getProducerId(), send.getLowestSequenceId(), send.getHighestSequenceId(), + if (send.hasLastSequenceId() && send.getSequenceId() <= send.getLastSequenceId()) { + producer.publishMessage(send.getProducerId(), send.getSequenceId(), send.getLastSequenceId(), headersAndPayload, send.getNumMessages()); } else { producer.publishMessage(send.getProducerId(), send.getSequenceId(), headersAndPayload, send.getNumMessages()); 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 88143df38bd77..7a5ce1986279a 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 @@ -160,8 +160,8 @@ public void readEntriesComplete(List entries, Object ctx) { MessageMetadata md = Commands.parseMessageMetadata(messageMetadataAndPayload); String producerName = md.getProducerName(); - long sequenceId = md.hasHighestSequenceId() && md.getHighestSequenceId() > 0 ? - md.getHighestSequenceId() : md.getSequenceId(); + long sequenceId = md.hasLastSequenceId() && md.getLastSequenceId() >= md.getSequenceId() ? + md.getLastSequenceId() : md.getSequenceId(); highestSequencedPushed.put(producerName, sequenceId); highestSequencedPersisted.put(producerName, sequenceId); @@ -293,8 +293,8 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade MessageMetadata md = Commands.parseMessageMetadata(headersAndPayload); producerName = md.getProducerName(); sequenceId = md.getSequenceId(); - lowestSequenceId = md.getLowestSequenceId(); - highestSequenceId = md.getHighestSequenceId(); + lowestSequenceId = md.getSequenceId(); + highestSequenceId = md.getLastSequenceId(); publishContext.setOriginalProducerName(producerName); publishContext.setOriginalSequenceId(sequenceId); headersAndPayload.readerIndex(readerIndex); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchMessageContainerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchMessageContainerImpl.java index ad911d228c80e..7e5fe9efd1cb4 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchMessageContainerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchMessageContainerImpl.java @@ -83,7 +83,7 @@ public void add(MessageImpl msg, SendCallback callback) { messages.add(msg); if (lowestSequenceId == -1) { lowestSequenceId = msg.getSequenceId(); - messageMetadata.setLowestSequenceId(lowestSequenceId); + messageMetadata.setSequenceId(lowestSequenceId); } highestSequenceId = msg.getSequenceId(); producer.lastSequenceIdPushed = msg.getSequenceId(); @@ -172,12 +172,12 @@ public boolean isMultiBatches() { public OpSendMsg createOpSendMsg() throws IOException { ByteBuf encryptedPayload = producer.encryptMessage(messageMetadata, getCompressedBatchMetadataAndPayload()); messageMetadata.setNumMessagesInBatch(numMessagesInBatch); - messageMetadata.setHighestSequenceId(highestSequenceId); - ByteBufPair cmd = producer.sendMessage(producer.producerId, messageMetadata.getLowestSequenceId(), - messageMetadata.getHighestSequenceId(), numMessagesInBatch, messageMetadata.build(), encryptedPayload); + messageMetadata.setLastSequenceId(highestSequenceId); + ByteBufPair cmd = producer.sendMessage(producer.producerId, messageMetadata.getSequenceId(), + messageMetadata.getLastSequenceId(), numMessagesInBatch, messageMetadata.build(), encryptedPayload); - OpSendMsg op = OpSendMsg.create(messages, cmd, messageMetadata.getLowestSequenceId(), - messageMetadata.getHighestSequenceId(), firstCallback); + OpSendMsg op = OpSendMsg.create(messages, cmd, messageMetadata.getSequenceId(), + messageMetadata.getLastSequenceId(), firstCallback); if (encryptedPayload.readableBytes() > ClientCnx.getMaxMessageSize()) { cmd.release(); 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 00d9ed5d2a1c1..129f97930a231 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 @@ -371,18 +371,10 @@ public void sendAsync(Message message, SendCallback callback) { msgMetadataBuilder.setSequenceId(sequenceId); } else { sequenceId = msgMetadataBuilder.getSequenceId(); - } - - if (sequenceId <= lastSequenceIdPushed) { - if (sequenceId <= lastSequenceIdPublished) { + if (sequenceId <= lastSequenceIdPushed) { callback.sendComplete(new PulsarClientException .InvalidMessageException("Message is definitely a duplicate")); return; - } else { - callback.sendComplete(new PulsarClientException - .InvalidMessageException("Message is a definitely a duplicate or not cannot be " + - "determined at this time")); - return; } } if (!msgMetadataBuilder.hasPublishTime()) { diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java b/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java index 145051ace9654..dd0c587d9e74b 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java @@ -3590,13 +3590,9 @@ public interface MessageMetadataOrBuilder boolean hasTxnidMostBits(); long getTxnidMostBits(); - // optional uint64 lowest_sequence_id = 24 [default = 0]; - boolean hasLowestSequenceId(); - long getLowestSequenceId(); - - // optional uint64 highest_sequence_id = 25 [default = 0]; - boolean hasHighestSequenceId(); - long getHighestSequenceId(); + // optional uint64 last_sequence_id = 24 [default = 0]; + boolean hasLastSequenceId(); + long getLastSequenceId(); } public static final class MessageMetadata extends org.apache.pulsar.shaded.com.google.protobuf.v241.GeneratedMessageLite @@ -3957,24 +3953,14 @@ public long getTxnidMostBits() { return txnidMostBits_; } - // optional uint64 lowest_sequence_id = 24 [default = 0]; - public static final int LOWEST_SEQUENCE_ID_FIELD_NUMBER = 24; - private long lowestSequenceId_; - public boolean hasLowestSequenceId() { + // optional uint64 last_sequence_id = 24 [default = 0]; + public static final int LAST_SEQUENCE_ID_FIELD_NUMBER = 24; + private long lastSequenceId_; + public boolean hasLastSequenceId() { return ((bitField0_ & 0x00040000) == 0x00040000); } - public long getLowestSequenceId() { - return lowestSequenceId_; - } - - // optional uint64 highest_sequence_id = 25 [default = 0]; - public static final int HIGHEST_SEQUENCE_ID_FIELD_NUMBER = 25; - private long highestSequenceId_; - public boolean hasHighestSequenceId() { - return ((bitField0_ & 0x00080000) == 0x00080000); - } - public long getHighestSequenceId() { - return highestSequenceId_; + public long getLastSequenceId() { + return lastSequenceId_; } private void initFields() { @@ -3999,8 +3985,7 @@ private void initFields() { markerType_ = 0; txnidLeastBits_ = 0L; txnidMostBits_ = 0L; - lowestSequenceId_ = 0L; - highestSequenceId_ = 0L; + lastSequenceId_ = 0L; } private byte memoizedIsInitialized = -1; public final boolean isInitialized() { @@ -4107,10 +4092,7 @@ public void writeTo(org.apache.pulsar.common.util.protobuf.ByteBufCodedOutputStr output.writeUInt64(23, txnidMostBits_); } if (((bitField0_ & 0x00040000) == 0x00040000)) { - output.writeUInt64(24, lowestSequenceId_); - } - if (((bitField0_ & 0x00080000) == 0x00080000)) { - output.writeUInt64(25, highestSequenceId_); + output.writeUInt64(24, lastSequenceId_); } } @@ -4211,11 +4193,7 @@ public int getSerializedSize() { } if (((bitField0_ & 0x00040000) == 0x00040000)) { size += org.apache.pulsar.shaded.com.google.protobuf.v241.CodedOutputStream - .computeUInt64Size(24, lowestSequenceId_); - } - if (((bitField0_ & 0x00080000) == 0x00080000)) { - size += org.apache.pulsar.shaded.com.google.protobuf.v241.CodedOutputStream - .computeUInt64Size(25, highestSequenceId_); + .computeUInt64Size(24, lastSequenceId_); } memoizedSerializedSize = size; return size; @@ -4372,10 +4350,8 @@ public Builder clear() { bitField0_ = (bitField0_ & ~0x00080000); txnidMostBits_ = 0L; bitField0_ = (bitField0_ & ~0x00100000); - lowestSequenceId_ = 0L; + lastSequenceId_ = 0L; bitField0_ = (bitField0_ & ~0x00200000); - highestSequenceId_ = 0L; - bitField0_ = (bitField0_ & ~0x00400000); return this; } @@ -4500,11 +4476,7 @@ public org.apache.pulsar.common.api.proto.PulsarApi.MessageMetadata buildPartial if (((from_bitField0_ & 0x00200000) == 0x00200000)) { to_bitField0_ |= 0x00040000; } - result.lowestSequenceId_ = lowestSequenceId_; - if (((from_bitField0_ & 0x00400000) == 0x00400000)) { - to_bitField0_ |= 0x00080000; - } - result.highestSequenceId_ = highestSequenceId_; + result.lastSequenceId_ = lastSequenceId_; result.bitField0_ = to_bitField0_; return result; } @@ -4595,11 +4567,8 @@ public Builder mergeFrom(org.apache.pulsar.common.api.proto.PulsarApi.MessageMet if (other.hasTxnidMostBits()) { setTxnidMostBits(other.getTxnidMostBits()); } - if (other.hasLowestSequenceId()) { - setLowestSequenceId(other.getLowestSequenceId()); - } - if (other.hasHighestSequenceId()) { - setHighestSequenceId(other.getHighestSequenceId()); + if (other.hasLastSequenceId()) { + setLastSequenceId(other.getLastSequenceId()); } return this; } @@ -4767,12 +4736,7 @@ public Builder mergeFrom( } case 192: { bitField0_ |= 0x00200000; - lowestSequenceId_ = input.readUInt64(); - break; - } - case 200: { - bitField0_ |= 0x00400000; - highestSequenceId_ = input.readUInt64(); + lastSequenceId_ = input.readUInt64(); break; } } @@ -5465,44 +5429,23 @@ public Builder clearTxnidMostBits() { return this; } - // optional uint64 lowest_sequence_id = 24 [default = 0]; - private long lowestSequenceId_ ; - public boolean hasLowestSequenceId() { + // optional uint64 last_sequence_id = 24 [default = 0]; + private long lastSequenceId_ ; + public boolean hasLastSequenceId() { return ((bitField0_ & 0x00200000) == 0x00200000); } - public long getLowestSequenceId() { - return lowestSequenceId_; + public long getLastSequenceId() { + return lastSequenceId_; } - public Builder setLowestSequenceId(long value) { + public Builder setLastSequenceId(long value) { bitField0_ |= 0x00200000; - lowestSequenceId_ = value; + lastSequenceId_ = value; return this; } - public Builder clearLowestSequenceId() { + public Builder clearLastSequenceId() { bitField0_ = (bitField0_ & ~0x00200000); - lowestSequenceId_ = 0L; - - return this; - } - - // optional uint64 highest_sequence_id = 25 [default = 0]; - private long highestSequenceId_ ; - public boolean hasHighestSequenceId() { - return ((bitField0_ & 0x00400000) == 0x00400000); - } - public long getHighestSequenceId() { - return highestSequenceId_; - } - public Builder setHighestSequenceId(long value) { - bitField0_ |= 0x00400000; - highestSequenceId_ = value; - - return this; - } - public Builder clearHighestSequenceId() { - bitField0_ = (bitField0_ & ~0x00400000); - highestSequenceId_ = 0L; + lastSequenceId_ = 0L; return this; } @@ -15442,13 +15385,9 @@ public interface CommandSendOrBuilder boolean hasTxnidMostBits(); long getTxnidMostBits(); - // optional uint64 lowest_sequence_id = 6 [default = 0]; - boolean hasLowestSequenceId(); - long getLowestSequenceId(); - - // optional uint64 highest_sequence_id = 7 [default = 0]; - boolean hasHighestSequenceId(); - long getHighestSequenceId(); + // optional uint64 last_sequence_id = 6 [default = 0]; + boolean hasLastSequenceId(); + long getLastSequenceId(); } public static final class CommandSend extends org.apache.pulsar.shaded.com.google.protobuf.v241.GeneratedMessageLite @@ -15535,24 +15474,14 @@ public long getTxnidMostBits() { return txnidMostBits_; } - // optional uint64 lowest_sequence_id = 6 [default = 0]; - public static final int LOWEST_SEQUENCE_ID_FIELD_NUMBER = 6; - private long lowestSequenceId_; - public boolean hasLowestSequenceId() { + // optional uint64 last_sequence_id = 6 [default = 0]; + public static final int LAST_SEQUENCE_ID_FIELD_NUMBER = 6; + private long lastSequenceId_; + public boolean hasLastSequenceId() { return ((bitField0_ & 0x00000020) == 0x00000020); } - public long getLowestSequenceId() { - return lowestSequenceId_; - } - - // optional uint64 highest_sequence_id = 7 [default = 0]; - public static final int HIGHEST_SEQUENCE_ID_FIELD_NUMBER = 7; - private long highestSequenceId_; - public boolean hasHighestSequenceId() { - return ((bitField0_ & 0x00000040) == 0x00000040); - } - public long getHighestSequenceId() { - return highestSequenceId_; + public long getLastSequenceId() { + return lastSequenceId_; } private void initFields() { @@ -15561,8 +15490,7 @@ private void initFields() { numMessages_ = 1; txnidLeastBits_ = 0L; txnidMostBits_ = 0L; - lowestSequenceId_ = 0L; - highestSequenceId_ = 0L; + lastSequenceId_ = 0L; } private byte memoizedIsInitialized = -1; public final boolean isInitialized() { @@ -15605,10 +15533,7 @@ public void writeTo(org.apache.pulsar.common.util.protobuf.ByteBufCodedOutputStr output.writeUInt64(5, txnidMostBits_); } if (((bitField0_ & 0x00000020) == 0x00000020)) { - output.writeUInt64(6, lowestSequenceId_); - } - if (((bitField0_ & 0x00000040) == 0x00000040)) { - output.writeUInt64(7, highestSequenceId_); + output.writeUInt64(6, lastSequenceId_); } } @@ -15640,11 +15565,7 @@ public int getSerializedSize() { } if (((bitField0_ & 0x00000020) == 0x00000020)) { size += org.apache.pulsar.shaded.com.google.protobuf.v241.CodedOutputStream - .computeUInt64Size(6, lowestSequenceId_); - } - if (((bitField0_ & 0x00000040) == 0x00000040)) { - size += org.apache.pulsar.shaded.com.google.protobuf.v241.CodedOutputStream - .computeUInt64Size(7, highestSequenceId_); + .computeUInt64Size(6, lastSequenceId_); } memoizedSerializedSize = size; return size; @@ -15769,10 +15690,8 @@ public Builder clear() { bitField0_ = (bitField0_ & ~0x00000008); txnidMostBits_ = 0L; bitField0_ = (bitField0_ & ~0x00000010); - lowestSequenceId_ = 0L; + lastSequenceId_ = 0L; bitField0_ = (bitField0_ & ~0x00000020); - highestSequenceId_ = 0L; - bitField0_ = (bitField0_ & ~0x00000040); return this; } @@ -15829,11 +15748,7 @@ public org.apache.pulsar.common.api.proto.PulsarApi.CommandSend buildPartial() { if (((from_bitField0_ & 0x00000020) == 0x00000020)) { to_bitField0_ |= 0x00000020; } - result.lowestSequenceId_ = lowestSequenceId_; - if (((from_bitField0_ & 0x00000040) == 0x00000040)) { - to_bitField0_ |= 0x00000040; - } - result.highestSequenceId_ = highestSequenceId_; + result.lastSequenceId_ = lastSequenceId_; result.bitField0_ = to_bitField0_; return result; } @@ -15855,11 +15770,8 @@ public Builder mergeFrom(org.apache.pulsar.common.api.proto.PulsarApi.CommandSen if (other.hasTxnidMostBits()) { setTxnidMostBits(other.getTxnidMostBits()); } - if (other.hasLowestSequenceId()) { - setLowestSequenceId(other.getLowestSequenceId()); - } - if (other.hasHighestSequenceId()) { - setHighestSequenceId(other.getHighestSequenceId()); + if (other.hasLastSequenceId()) { + setLastSequenceId(other.getLastSequenceId()); } return this; } @@ -15925,12 +15837,7 @@ public Builder mergeFrom( } case 48: { bitField0_ |= 0x00000020; - lowestSequenceId_ = input.readUInt64(); - break; - } - case 56: { - bitField0_ |= 0x00000040; - highestSequenceId_ = input.readUInt64(); + lastSequenceId_ = input.readUInt64(); break; } } @@ -16044,44 +15951,23 @@ public Builder clearTxnidMostBits() { return this; } - // optional uint64 lowest_sequence_id = 6 [default = 0]; - private long lowestSequenceId_ ; - public boolean hasLowestSequenceId() { + // optional uint64 last_sequence_id = 6 [default = 0]; + private long lastSequenceId_ ; + public boolean hasLastSequenceId() { return ((bitField0_ & 0x00000020) == 0x00000020); } - public long getLowestSequenceId() { - return lowestSequenceId_; + public long getLastSequenceId() { + return lastSequenceId_; } - public Builder setLowestSequenceId(long value) { + public Builder setLastSequenceId(long value) { bitField0_ |= 0x00000020; - lowestSequenceId_ = value; + lastSequenceId_ = value; return this; } - public Builder clearLowestSequenceId() { + public Builder clearLastSequenceId() { bitField0_ = (bitField0_ & ~0x00000020); - lowestSequenceId_ = 0L; - - return this; - } - - // optional uint64 highest_sequence_id = 7 [default = 0]; - private long highestSequenceId_ ; - public boolean hasHighestSequenceId() { - return ((bitField0_ & 0x00000040) == 0x00000040); - } - public long getHighestSequenceId() { - return highestSequenceId_; - } - public Builder setHighestSequenceId(long value) { - bitField0_ |= 0x00000040; - highestSequenceId_ = value; - - return this; - } - public Builder clearHighestSequenceId() { - bitField0_ = (bitField0_ & ~0x00000040); - highestSequenceId_ = 0L; + lastSequenceId_ = 0L; return this; } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java index 00c35307e5f9d..48e70d8f93c0e 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java @@ -478,8 +478,8 @@ public static ByteBufPair newSend(long producerId, long lowestSequenceId, long h MessageMetadata messageData, ByteBuf payload) { CommandSend.Builder sendBuilder = CommandSend.newBuilder(); sendBuilder.setProducerId(producerId); - sendBuilder.setLowestSequenceId(lowestSequenceId); - sendBuilder.setHighestSequenceId(highestSequenceId); + sendBuilder.setSequenceId(lowestSequenceId); + sendBuilder.setLastSequenceId(highestSequenceId); sendBuilder.setSequenceId(lowestSequenceId); if (numMessages > 1) { sendBuilder.setNumMessages(numMessages); @@ -1547,7 +1547,6 @@ public static long initBatchMessageMetadata(MessageMetadata.Builder messageMetad messageMetadata.setPublishTime(builder.getPublishTime()); messageMetadata.setProducerName(builder.getProducerName()); messageMetadata.setSequenceId(builder.getSequenceId()); - messageMetadata.setLowestSequenceId(builder.getSequenceId()); if (builder.hasReplicatedFrom()) { messageMetadata.setReplicatedFrom(builder.getReplicatedFrom()); } diff --git a/pulsar-common/src/main/proto/PulsarApi.proto b/pulsar-common/src/main/proto/PulsarApi.proto index d76d2c4a33048..9b99f3539f7d0 100644 --- a/pulsar-common/src/main/proto/PulsarApi.proto +++ b/pulsar-common/src/main/proto/PulsarApi.proto @@ -133,9 +133,8 @@ message MessageMetadata { optional uint64 txnid_least_bits = 22 [default = 0]; optional uint64 txnid_most_bits = 23 [default = 0]; - /// Add lowest and highest sequence id to support external sequence id - optional uint64 lowest_sequence_id = 24 [default = 0]; - optional uint64 highest_sequence_id = 25 [default = 0]; + /// Add last sequence id to support batch message with external sequence id + optional uint64 last_sequence_id = 24 [default = 0]; } message SingleMessageMetadata { @@ -416,9 +415,8 @@ message CommandSend { optional uint64 txnid_least_bits = 4 [default = 0]; optional uint64 txnid_most_bits = 5 [default = 0]; - /// Add lowest and highest sequence id to support external sequence id - optional uint64 lowest_sequence_id = 6 [default = 0]; - optional uint64 highest_sequence_id = 7 [default = 0]; + /// Add last sequence id to support batch message with external sequence id + optional uint64 last_sequence_id = 6 [default = 0]; } message CommandSendReceipt { From 126522cc0199a1db4e4aa611efc5f0e08b283aed Mon Sep 17 00:00:00 2001 From: lipenghui Date: Fri, 1 Nov 2019 11:28:34 +0800 Subject: [PATCH 08/14] fix comments --- .../pulsar/broker/service/Producer.java | 51 +++++++++++-------- .../apache/pulsar/broker/service/Topic.java | 10 ++-- .../persistent/MessageDeduplication.java | 32 ++++-------- .../pulsar/client/impl/ProducerImpl.java | 9 +++- 4 files changed, 53 insertions(+), 49 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java index 3ee6664b384f0..d5fbfb42eefbb 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java @@ -136,18 +136,18 @@ public void publishMessage(long producerId, long sequenceId, ByteBuf headersAndP publishMessageToTopic(headersAndPayload, sequenceId, batchSize); } - public void publishMessage(long producerId, long lowestSequenceId, long highestSequenceId, + public void publishMessage(long producerId, long sequenceId, long lastSequenceId, ByteBuf headersAndPayload, long batchSize) { - if (lowestSequenceId > highestSequenceId) { + if (sequenceId > lastSequenceId) { cnx.ctx().channel().eventLoop().execute(() -> { - cnx.ctx().writeAndFlush(Commands.newSendError(producerId, highestSequenceId, ServerError.MetadataError, + cnx.ctx().writeAndFlush(Commands.newSendError(producerId, lastSequenceId, ServerError.MetadataError, "Invalid lowest or highest sequence id")); cnx.completedSendOperation(isNonPersistentTopic); }); return; } - beforePublish(producerId, highestSequenceId, headersAndPayload, batchSize); - publishMessageToTopic(headersAndPayload, lowestSequenceId, highestSequenceId, batchSize); + beforePublish(producerId, lastSequenceId, headersAndPayload, batchSize); + publishMessageToTopic(headersAndPayload, sequenceId, lastSequenceId, batchSize); } public void beforePublish(long producerId, long sequenceId, ByteBuf headersAndPayload, long batchSize) { @@ -197,9 +197,9 @@ private void publishMessageToTopic(ByteBuf headersAndPayload, long sequenceId, l System.nanoTime())); } - private void publishMessageToTopic(ByteBuf headersAndPayload, long lowestSequenceId, long highestSequenceId, long batchSize) { + private void publishMessageToTopic(ByteBuf headersAndPayload, long sequenceId, long lastSequenceId, long batchSize) { topic.publishMessage(headersAndPayload, - MessagePublishContext.get(this, lowestSequenceId, highestSequenceId, msgIn, headersAndPayload.readableBytes(), batchSize, + MessagePublishContext.get(this, sequenceId, lastSequenceId, msgIn, headersAndPayload.readableBytes(), batchSize, System.nanoTime())); } @@ -285,8 +285,8 @@ private static final class MessagePublishContext implements PublishContext, Runn private String originalProducerName; private long originalSequenceId; - private long lowestSequenceId; - private long highestSequenceId; + private long lastSequenceId; + private long originalLastSequenceId; public String getProducerName() { return producer.getProducerName(); @@ -296,13 +296,9 @@ public long getSequenceId() { return sequenceId; } - public long getLowestSequenceId() { - return lowestSequenceId; - } - @Override - public long getHighestSequenceId() { - return highestSequenceId; + public long getLastSequenceId() { + return lastSequenceId; } @Override @@ -325,6 +321,16 @@ public long getOriginalSequenceId() { return originalSequenceId; } + @Override + public void setOriginalLastSequenceId(long originalLastSequenceId) { + this.originalLastSequenceId = originalLastSequenceId; + } + + @Override + public long getOriginalLastSequenceId() { + return originalLastSequenceId; + } + /** * Executed from managed ledger thread when the message is persisted */ @@ -338,7 +344,7 @@ public void completed(Exception exception, long ledgerId, long entryId) { if (!(exception instanceof TopicClosedException)) { // For TopicClosed exception there's no need to send explicit error, since the client was // already notified - long callBackSequenceId = highestSequenceId > 0 ? highestSequenceId : sequenceId; + long callBackSequenceId = lastSequenceId >= sequenceId ? lastSequenceId : sequenceId; producer.cnx.ctx().writeAndFlush(Commands.newSendError(producer.producerId, callBackSequenceId, serverError, exception.getMessage())); } @@ -371,7 +377,7 @@ public void run() { // stats rateIn.recordMultipleEvents(batchSize, msgSize); producer.topic.recordAddLatency(System.nanoTime() - startTimeNs, TimeUnit.NANOSECONDS); - long callBackSequenceId = highestSequenceId >= 0 && highestSequenceId >= sequenceId ? highestSequenceId : sequenceId; + long callBackSequenceId = lastSequenceId >= sequenceId ? lastSequenceId : sequenceId;; producer.cnx.ctx().writeAndFlush( Commands.newSendReceipt(producer.producerId, callBackSequenceId, ledgerId, entryId), producer.cnx.ctx().voidPromise()); @@ -394,12 +400,12 @@ static MessagePublishContext get(Producer producer, long sequenceId, Rate rateIn return callback; } - static MessagePublishContext get(Producer producer, long lowestSequenceId, long highestSequenceId, Rate rateIn, + static MessagePublishContext get(Producer producer, long sequenceId, long lastSequenceId, Rate rateIn, int msgSize, long batchSize, long startTimeNs) { MessagePublishContext callback = RECYCLER.get(); callback.producer = producer; - callback.lowestSequenceId = lowestSequenceId; - callback.highestSequenceId = highestSequenceId; + callback.sequenceId = sequenceId; + callback.lastSequenceId = lastSequenceId; callback.rateIn = rateIn; callback.msgSize = msgSize; callback.batchSize = batchSize; @@ -424,8 +430,9 @@ protected MessagePublishContext newObject(Recycler.Handle public void recycle() { producer = null; sequenceId = -1; - lowestSequenceId = -1; - highestSequenceId = -1; + lastSequenceId = -1; + originalSequenceId = -1; + originalLastSequenceId = -1; rateIn = null; msgSize = 0; ledgerId = -1; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java index e668562bbd186..ccc871e484eb7 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java @@ -78,11 +78,15 @@ default long getOriginalSequenceId() { void completed(Exception e, long ledgerId, long entryId); - default long getLowestSequenceId() { - return -1; + default long getLastSequenceId() { + return -1; + } + + default void setOriginalLastSequenceId(long originalLastSequenceId) { + } - default long getHighestSequenceId() { + default long getOriginalLastSequenceId() { return -1; } } 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 7a5ce1986279a..f1234c4955783 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 @@ -284,8 +284,8 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade String producerName = publishContext.getProducerName(); long sequenceId = publishContext.getSequenceId(); - long lowestSequenceId = publishContext.getLowestSequenceId(); - long highestSequenceId = publishContext.getHighestSequenceId(); + long lastSequenceId = publishContext.getLastSequenceId() >= sequenceId ? publishContext.getLastSequenceId() + : 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. @@ -293,24 +293,19 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade MessageMetadata md = Commands.parseMessageMetadata(headersAndPayload); producerName = md.getProducerName(); sequenceId = md.getSequenceId(); - lowestSequenceId = md.getSequenceId(); - highestSequenceId = md.getLastSequenceId(); + lastSequenceId = md.getLastSequenceId() >= sequenceId ? md.getLastSequenceId() : md.getSequenceId(); publishContext.setOriginalProducerName(producerName); publishContext.setOriginalSequenceId(sequenceId); + publishContext.setOriginalLastSequenceId(lastSequenceId); headersAndPayload.readerIndex(readerIndex); md.recycle(); } - if (lowestSequenceId == highestSequenceId && lowestSequenceId <= sequenceId) { - lowestSequenceId = sequenceId; - highestSequenceId = sequenceId; - } - // 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 && lowestSequenceId <= 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); @@ -321,14 +316,13 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade // If current message's seq id is between lastSequenceIdPersisted and 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 && lowestSequenceId <= lastSequenceIdPersisted) { + if (lastSequenceIdPersisted != null && sequenceId <= lastSequenceIdPersisted) { return MessageDupStatus.Dup; } else { return MessageDupStatus.Unknown; } } - - highestSequencedPushed.put(producerName, highestSequenceId); + highestSequencedPushed.put(producerName, lastSequenceId); } return MessageDupStatus.NotDup; } @@ -343,21 +337,15 @@ public void recordMessagePersisted(PublishContext publishContext, PositionImpl p String producerName = publishContext.getProducerName(); long sequenceId = publishContext.getSequenceId(); - long lowestSequenceId = publishContext.getLowestSequenceId(); - long highestSequenceId = publishContext.getHighestSequenceId(); + long lastSequenceId = publishContext.getLastSequenceId(); if (publishContext.getOriginalProducerName() != null) { // In case of replicated messages, this will be different from the current replicator producer name producerName = publishContext.getOriginalProducerName(); sequenceId = publishContext.getOriginalSequenceId(); - lowestSequenceId = publishContext.getLowestSequenceId(); - highestSequenceId = publishContext.getHighestSequenceId(); - } - - if (lowestSequenceId <= highestSequenceId && lowestSequenceId <= 0) { - highestSequenceId = sequenceId; + lastSequenceId = publishContext.getOriginalLastSequenceId(); } - highestSequencedPersisted.put(producerName, highestSequenceId); + highestSequencedPersisted.put(producerName, lastSequenceId >= sequenceId ? lastSequenceId : sequenceId); if (++snapshotCounter >= snapshotInterval) { snapshotCounter = 0; takeSnapshot(position); 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 129f97930a231..045c3ba2d401e 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 @@ -372,8 +372,13 @@ public void sendAsync(Message message, SendCallback callback) { } else { sequenceId = msgMetadataBuilder.getSequenceId(); if (sequenceId <= lastSequenceIdPushed) { - callback.sendComplete(new PulsarClientException - .InvalidMessageException("Message is definitely a duplicate")); + if (sequenceId <= lastSequenceIdPublished) { + log.warn("Message with sequence id {} is definitely a duplicate", sequenceId); + } else { + log.warn("Message with sequence id {} is a definitely a duplicate or not cannot be " + + "determined at this time", sequenceId); + } + callback.getFuture().complete(null); return; } } From b4d158d5c2dd8ead9999c515dc09be47be8609f5 Mon Sep 17 00:00:00 2001 From: lipenghui Date: Fri, 1 Nov 2019 14:16:44 +0800 Subject: [PATCH 09/14] fix comments --- .../broker/service/BatchMessageTest.java | 2 - .../persistent/MessageDuplicationTest.java | 11 +++--- .../api/ClientDeduplicationFailureTest.java | 6 +-- .../client/api/ClientDeduplicationTest.java | 20 +++------- .../pulsar/client/api/SimpleSchemaTest.java | 2 +- .../apache/pulsar/client/impl/ReaderTest.java | 2 +- .../pulsar/client/impl/ProducerImpl.java | 38 ++++++++++--------- 7 files changed, 35 insertions(+), 46 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageTest.java index f3e024909a769..b0ca4b4368d34 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageTest.java @@ -567,14 +567,12 @@ public void testBatchAndNonBatchCumulativeAcks(BatcherBuilder builder) throws Ex .batchingMaxMessages(numMsgsInBatch) .enableBatching(true) .batcherBuilder(builder) - .producerName("1") .messageRoutingMode(MessageRoutingMode.SinglePartition) .create(); // create producer to publish non batch messages Producer noBatchProducer = pulsarClient.newProducer().topic(topicName) .enableBatching(false) .messageRoutingMode(MessageRoutingMode.SinglePartition) - .producerName("2") .create(); List> sendFutureList = Lists.newArrayList(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java index 0b205ddf505e0..a683b05fabe05 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java @@ -339,20 +339,21 @@ public void completed(Exception e, long ledgerId, long entryId) { }); } - public Topic.PublishContext getPublishContext(String producerName, long lowestSequenceId, long highestSequenceId) { + public Topic.PublishContext getPublishContext(String producerName, long seqId, long lastSequenceId) { return spy(new Topic.PublishContext() { @Override public String getProducerName() { return producerName; } - public long getLowestSequenceId() { - return lowestSequenceId; + @Override + public long getSequenceId() { + return seqId; } @Override - public long getHighestSequenceId() { - return highestSequenceId; + public long getLastSequenceId() { + return lastSequenceId; } @Override diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ClientDeduplicationFailureTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ClientDeduplicationFailureTest.java index 567e9db87973f..63f8aaa2bcd45 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ClientDeduplicationFailureTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ClientDeduplicationFailureTest.java @@ -241,11 +241,7 @@ public void testClientDeduplicationCorrectnessWithFailure() throws Exception { producerThread.stop(); // send last message - try { - producer.newMessage().sequenceId(producerThread.getLastSeqId() + 1).value("end").send(); - fail("should failed, because send a duplication"); - } catch (PulsarClientException.InvalidMessageException ignore) { - } + producer.newMessage().sequenceId(producerThread.getLastSeqId() + 1).value("end").send(); producer.close(); Reader reader = pulsarClient.newReader(Schema.STRING).startMessageId(MessageId.earliest).topic(sourceTopic).create(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ClientDeduplicationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ClientDeduplicationTest.java index 916873ebeb0a0..97e7869a8bb33 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ClientDeduplicationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ClientDeduplicationTest.java @@ -129,12 +129,8 @@ public void testProducerDeduplication() throws Exception { producer.newMessage().value("my-message-2".getBytes()).sequenceId(2).send(); // Repeat the messages and verify they're not received by consumer - try { - producer.newMessage().value("my-message-1".getBytes()).sequenceId(1).send(); - producer.newMessage().value("my-message-2".getBytes()).sequenceId(2).send(); - fail("should be failed"); - } catch (PulsarClientException ignore) { - } + producer.newMessage().value("my-message-1".getBytes()).sequenceId(1).send(); + producer.newMessage().value("my-message-2".getBytes()).sequenceId(2).send(); producer.close(); @@ -186,10 +182,8 @@ public void testProducerDeduplicationWithDiscontinuousSequenceId() throws Except producer.flush(); // Repeat the messages and verify they're not received by consumer - CompletableFuture sendResult = producer.newMessage().value("my-message-1".getBytes()).sequenceId(2).sendAsync(); - assertTrue(sendResult.isCompletedExceptionally()); - sendResult = producer.newMessage().value("my-message-2".getBytes()).sequenceId(4).sendAsync(); - assertTrue(sendResult.isCompletedExceptionally()); + producer.newMessage().value("my-message-1".getBytes()).sequenceId(2).sendAsync(); + producer.newMessage().value("my-message-2".getBytes()).sequenceId(4).sendAsync(); producer.close(); for (int i = 0; i < 3; i++) { @@ -239,10 +233,8 @@ public void testProducerDeduplicationNonBatchAsync() throws Exception { producer.newMessage().value("my-message-2".getBytes()).sequenceId(5).sendAsync(); // Repeat the messages and verify they're not received by consumer - CompletableFuture sendResult = producer.newMessage().value("my-message-1".getBytes()).sequenceId(2).sendAsync(); - assertTrue(sendResult.isCompletedExceptionally()); - sendResult = producer.newMessage().value("my-message-2".getBytes()).sequenceId(4).sendAsync(); - assertTrue(sendResult.isCompletedExceptionally()); + producer.newMessage().value("my-message-1".getBytes()).sequenceId(2).sendAsync(); + producer.newMessage().value("my-message-2".getBytes()).sequenceId(4).sendAsync(); producer.close(); for (int i = 0; i < 3; i++) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleSchemaTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleSchemaTest.java index 5007b2f098444..aa0f0e0c9018b 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleSchemaTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleSchemaTest.java @@ -301,7 +301,7 @@ public void newProducerForMessageOnTopicWithDifferentSchemaType() throws Excepti @Test public void newProducerForMessageSchemaOnTopicInitialWithNoSchema() throws Exception { - String topic = "my-property/my-ns/schema-test-"; + String topic = "my-property/my-ns/schema-test"; Schema v1Schema = Schema.AVRO(V1Data.class); byte[] v1SchemaBytes = v1Schema.getSchemaInfo().getSchema(); AvroWriter v1Writer = new AvroWriter<>( diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/ReaderTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/ReaderTest.java index 2919860855a8d..e931c6c2fa017 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/ReaderTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/ReaderTest.java @@ -189,7 +189,7 @@ public void testReaderWithTimeLong() throws Exception { TypedMessageBuilderImpl msg = (TypedMessageBuilderImpl) producer.newMessage() .value(("new" + i).getBytes()); Builder metadataBuilder = msg.getMetadataBuilder(); - metadataBuilder.setPublishTime(newMsgPublishTime).setSequenceId(totalMsg + i); + metadataBuilder.setPublishTime(newMsgPublishTime); metadataBuilder.setProducerName(producer.getProducerName()).setReplicatedFrom("us-west1"); MessageId msgId = msg.send(); if (firstMsgId == null) { 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 045c3ba2d401e..3d5c58b219075 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 @@ -371,16 +371,6 @@ public void sendAsync(Message message, SendCallback callback) { msgMetadataBuilder.setSequenceId(sequenceId); } else { sequenceId = msgMetadataBuilder.getSequenceId(); - if (sequenceId <= lastSequenceIdPushed) { - if (sequenceId <= lastSequenceIdPublished) { - log.warn("Message with sequence id {} is definitely a duplicate", sequenceId); - } else { - log.warn("Message with sequence id {} is a definitely a duplicate or not cannot be " + - "determined at this time", sequenceId); - } - callback.getFuture().complete(null); - return; - } } if (!msgMetadataBuilder.hasPublishTime()) { msgMetadataBuilder.setPublishTime(client.getClientClock().millis()); @@ -396,15 +386,27 @@ public void sendAsync(Message message, SendCallback callback) { msgMetadataBuilder.setUncompressedSize(uncompressedSize); } if (canAddToBatch(msg)) { - // handle boundary cases where message being added would exceed - // batch size and/or max message size if (canAddToCurrentBatch(msg)) { - batchMessageContainer.add(msg, callback); - lastSendFuture = callback.getFuture(); - payload.release(); - if (batchMessageContainer.getNumMessagesInBatch() == maxNumMessagesInBatch - || batchMessageContainer.getCurrentBatchSize() >= BatchMessageContainerImpl.MAX_MESSAGE_BATCH_SIZE_BYTES) { - batchMessageAndSend(); + // should trigger complete the batch message, new message will add to a new batch and new batch + // sequence id use the new message, so that broker can handle the message duplication + if (sequenceId <= lastSequenceIdPushed) { + if (sequenceId <= lastSequenceIdPublished) { + log.warn("Message with sequence id {} is definitely a duplicate", sequenceId); + } else { + log.warn("Message with sequence id {} is a definitely a duplicate or not cannot be " + + "determined at this time", sequenceId); + } + doBatchSendAndAdd(msg, callback, payload); + } else { + // handle boundary cases where message being added would exceed + // batch size and/or max message size + batchMessageContainer.add(msg, callback); + lastSendFuture = callback.getFuture(); + payload.release(); + if (batchMessageContainer.getNumMessagesInBatch() == maxNumMessagesInBatch + || batchMessageContainer.getCurrentBatchSize() >= BatchMessageContainerImpl.MAX_MESSAGE_BATCH_SIZE_BYTES) { + batchMessageAndSend(); + } } } else { doBatchSendAndAdd(msg, callback, payload); From 55540171eb2c28bf240ed40aedb8c8713f057771 Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Tue, 5 Nov 2019 20:05:51 +0800 Subject: [PATCH 10/14] Fix comments --- .../pulsar/broker/service/Producer.java | 36 +++--- .../pulsar/broker/service/ServerCnx.java | 4 +- .../apache/pulsar/broker/service/Topic.java | 14 +-- .../persistent/MessageDeduplication.java | 14 +-- .../persistent/MessageDuplicationTest.java | 2 +- .../impl/BatchMessageContainerImpl.java | 18 +-- .../pulsar/client/impl/ProducerImpl.java | 10 +- .../pulsar/common/api/proto/PulsarApi.java | 104 +++++++++--------- .../pulsar/common/protocol/Commands.java | 3 +- pulsar-common/src/main/proto/PulsarApi.proto | 8 +- 10 files changed, 108 insertions(+), 105 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java index d5fbfb42eefbb..d7366f9083be7 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java @@ -136,18 +136,18 @@ public void publishMessage(long producerId, long sequenceId, ByteBuf headersAndP publishMessageToTopic(headersAndPayload, sequenceId, batchSize); } - public void publishMessage(long producerId, long sequenceId, long lastSequenceId, + public void publishMessage(long producerId, long lowestSequenceId, long highestSequenceId, ByteBuf headersAndPayload, long batchSize) { - if (sequenceId > lastSequenceId) { + if (lowestSequenceId > highestSequenceId) { cnx.ctx().channel().eventLoop().execute(() -> { - cnx.ctx().writeAndFlush(Commands.newSendError(producerId, lastSequenceId, ServerError.MetadataError, + cnx.ctx().writeAndFlush(Commands.newSendError(producerId, highestSequenceId, ServerError.MetadataError, "Invalid lowest or highest sequence id")); cnx.completedSendOperation(isNonPersistentTopic); }); return; } - beforePublish(producerId, lastSequenceId, headersAndPayload, batchSize); - publishMessageToTopic(headersAndPayload, sequenceId, lastSequenceId, batchSize); + beforePublish(producerId, highestSequenceId, headersAndPayload, batchSize); + publishMessageToTopic(headersAndPayload, lowestSequenceId, highestSequenceId, batchSize); } public void beforePublish(long producerId, long sequenceId, ByteBuf headersAndPayload, long batchSize) { @@ -285,8 +285,8 @@ private static final class MessagePublishContext implements PublishContext, Runn private String originalProducerName; private long originalSequenceId; - private long lastSequenceId; - private long originalLastSequenceId; + private long highestSequenceId; + private long originalHighestSequenceId; public String getProducerName() { return producer.getProducerName(); @@ -297,8 +297,8 @@ public long getSequenceId() { } @Override - public long getLastSequenceId() { - return lastSequenceId; + public long getHighestSequenceId() { + return highestSequenceId; } @Override @@ -322,13 +322,13 @@ public long getOriginalSequenceId() { } @Override - public void setOriginalLastSequenceId(long originalLastSequenceId) { - this.originalLastSequenceId = originalLastSequenceId; + public void setOriginalHighestSequenceId(long originalHighestSequenceId) { + this.originalHighestSequenceId = originalHighestSequenceId; } @Override - public long getOriginalLastSequenceId() { - return originalLastSequenceId; + public long getOriginalHighestSequenceId() { + return originalHighestSequenceId; } /** @@ -344,7 +344,7 @@ public void completed(Exception exception, long ledgerId, long entryId) { if (!(exception instanceof TopicClosedException)) { // For TopicClosed exception there's no need to send explicit error, since the client was // already notified - long callBackSequenceId = lastSequenceId >= sequenceId ? lastSequenceId : sequenceId; + long callBackSequenceId = Math.max(highestSequenceId, sequenceId); producer.cnx.ctx().writeAndFlush(Commands.newSendError(producer.producerId, callBackSequenceId, serverError, exception.getMessage())); } @@ -377,7 +377,7 @@ public void run() { // stats rateIn.recordMultipleEvents(batchSize, msgSize); producer.topic.recordAddLatency(System.nanoTime() - startTimeNs, TimeUnit.NANOSECONDS); - long callBackSequenceId = lastSequenceId >= sequenceId ? lastSequenceId : sequenceId;; + long callBackSequenceId = highestSequenceId >= sequenceId ? highestSequenceId : sequenceId;; producer.cnx.ctx().writeAndFlush( Commands.newSendReceipt(producer.producerId, callBackSequenceId, ledgerId, entryId), producer.cnx.ctx().voidPromise()); @@ -405,7 +405,7 @@ static MessagePublishContext get(Producer producer, long sequenceId, long lastSe MessagePublishContext callback = RECYCLER.get(); callback.producer = producer; callback.sequenceId = sequenceId; - callback.lastSequenceId = lastSequenceId; + callback.highestSequenceId = lastSequenceId; callback.rateIn = rateIn; callback.msgSize = msgSize; callback.batchSize = batchSize; @@ -430,9 +430,9 @@ protected MessagePublishContext newObject(Recycler.Handle public void recycle() { producer = null; sequenceId = -1; - lastSequenceId = -1; + highestSequenceId = -1; originalSequenceId = -1; - originalLastSequenceId = -1; + originalHighestSequenceId = -1; rateIn = null; msgSize = 0; ledgerId = -1; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index 5ceb76a5b753e..6b38d377b9773 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -1051,8 +1051,8 @@ protected void handleSend(CommandSend send, ByteBuf headersAndPayload) { startSendOperation(producer); // Persist the message - if (send.hasLastSequenceId() && send.getSequenceId() <= send.getLastSequenceId()) { - producer.publishMessage(send.getProducerId(), send.getSequenceId(), send.getLastSequenceId(), + if (send.hasHighestSequenceId() && send.getSequenceId() <= send.getHighestSequenceId()) { + producer.publishMessage(send.getProducerId(), send.getSequenceId(), send.getHighestSequenceId(), headersAndPayload, send.getNumMessages()); } else { producer.publishMessage(send.getProducerId(), send.getSequenceId(), headersAndPayload, send.getNumMessages()); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java index ccc871e484eb7..d59b359d0d7a5 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java @@ -53,7 +53,7 @@ default String getProducerName() { } default long getSequenceId() { - return -1; + return -1L; } default void setOriginalProducerName(String originalProducerName) { @@ -73,21 +73,21 @@ default String getOriginalProducerName() { } default long getOriginalSequenceId() { - return -1; + return -1L; } void completed(Exception e, long ledgerId, long entryId); - default long getLastSequenceId() { - return -1; + default long getHighestSequenceId() { + return -1L; } - default void setOriginalLastSequenceId(long originalLastSequenceId) { + default void setOriginalHighestSequenceId(long originalHighestSequenceId) { } - default long getOriginalLastSequenceId() { - return -1; + default long getOriginalHighestSequenceId() { + return -1L; } } 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 f1234c4955783..2794989930c0c 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 @@ -160,8 +160,8 @@ public void readEntriesComplete(List entries, Object ctx) { MessageMetadata md = Commands.parseMessageMetadata(messageMetadataAndPayload); String producerName = md.getProducerName(); - long sequenceId = md.hasLastSequenceId() && md.getLastSequenceId() >= md.getSequenceId() ? - md.getLastSequenceId() : md.getSequenceId(); + long sequenceId = md.hasHighestSequenceId() && md.getHighestSequenceId() >= md.getSequenceId() ? + md.getHighestSequenceId() : md.getSequenceId(); highestSequencedPushed.put(producerName, sequenceId); highestSequencedPersisted.put(producerName, sequenceId); @@ -284,7 +284,7 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade String producerName = publishContext.getProducerName(); long sequenceId = publishContext.getSequenceId(); - long lastSequenceId = publishContext.getLastSequenceId() >= sequenceId ? publishContext.getLastSequenceId() + long lastSequenceId = publishContext.getHighestSequenceId() >= sequenceId ? publishContext.getHighestSequenceId() : sequenceId; if (producerName.startsWith(replicatorPrefix)) { // Message is coming from replication, we need to use the original producer name and sequence id @@ -293,10 +293,10 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade MessageMetadata md = Commands.parseMessageMetadata(headersAndPayload); producerName = md.getProducerName(); sequenceId = md.getSequenceId(); - lastSequenceId = md.getLastSequenceId() >= sequenceId ? md.getLastSequenceId() : md.getSequenceId(); + lastSequenceId = md.getHighestSequenceId() >= sequenceId ? md.getHighestSequenceId() : md.getSequenceId(); publishContext.setOriginalProducerName(producerName); publishContext.setOriginalSequenceId(sequenceId); - publishContext.setOriginalLastSequenceId(lastSequenceId); + publishContext.setOriginalHighestSequenceId(lastSequenceId); headersAndPayload.readerIndex(readerIndex); md.recycle(); } @@ -337,12 +337,12 @@ public void recordMessagePersisted(PublishContext publishContext, PositionImpl p String producerName = publishContext.getProducerName(); long sequenceId = publishContext.getSequenceId(); - long lastSequenceId = publishContext.getLastSequenceId(); + long lastSequenceId = publishContext.getHighestSequenceId(); if (publishContext.getOriginalProducerName() != null) { // In case of replicated messages, this will be different from the current replicator producer name producerName = publishContext.getOriginalProducerName(); sequenceId = publishContext.getOriginalSequenceId(); - lastSequenceId = publishContext.getOriginalLastSequenceId(); + lastSequenceId = publishContext.getOriginalHighestSequenceId(); } highestSequencedPersisted.put(producerName, lastSequenceId >= sequenceId ? lastSequenceId : sequenceId); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java index a683b05fabe05..5cfdef839db5a 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java @@ -352,7 +352,7 @@ public long getSequenceId() { } @Override - public long getLastSequenceId() { + public long getHighestSequenceId() { return lastSequenceId; } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchMessageContainerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchMessageContainerImpl.java index 7e5fe9efd1cb4..24eaf9d5198f9 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchMessageContainerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchMessageContainerImpl.java @@ -50,8 +50,8 @@ class BatchMessageContainerImpl extends AbstractBatchMessageContainer { private PulsarApi.MessageMetadata.Builder messageMetadata = PulsarApi.MessageMetadata.newBuilder(); // sequence id for this batch which will be persisted as a single entry by broker - private long lowestSequenceId = -1; - private long highestSequenceId = -1; + private long lowestSequenceId = -1L; + private long highestSequenceId = -1L; private ByteBuf batchedMessageMetadataAndPayload; private List> messages = Lists.newArrayList(); protected SendCallback previousCallback = null; @@ -81,7 +81,7 @@ public void add(MessageImpl msg, SendCallback callback) { previousCallback = callback; currentBatchSizeBytes += msg.getDataBuffer().readableBytes(); messages.add(msg); - if (lowestSequenceId == -1) { + if (lowestSequenceId == -1L) { lowestSequenceId = msg.getSequenceId(); messageMetadata.setSequenceId(lowestSequenceId); } @@ -138,8 +138,8 @@ public void clear() { messageMetadata.clear(); numMessagesInBatch = 0; currentBatchSizeBytes = 0; - lowestSequenceId = -1; - highestSequenceId = -1; + lowestSequenceId = -1L; + highestSequenceId = -1L; batchedMessageMetadataAndPayload = null; } @@ -172,12 +172,12 @@ public boolean isMultiBatches() { public OpSendMsg createOpSendMsg() throws IOException { ByteBuf encryptedPayload = producer.encryptMessage(messageMetadata, getCompressedBatchMetadataAndPayload()); messageMetadata.setNumMessagesInBatch(numMessagesInBatch); - messageMetadata.setLastSequenceId(highestSequenceId); + messageMetadata.setHighestSequenceId(highestSequenceId); ByteBufPair cmd = producer.sendMessage(producer.producerId, messageMetadata.getSequenceId(), - messageMetadata.getLastSequenceId(), numMessagesInBatch, messageMetadata.build(), encryptedPayload); + messageMetadata.getHighestSequenceId(), numMessagesInBatch, messageMetadata.build(), encryptedPayload); OpSendMsg op = OpSendMsg.create(messages, cmd, messageMetadata.getSequenceId(), - messageMetadata.getLastSequenceId(), firstCallback); + messageMetadata.getHighestSequenceId(), firstCallback); if (encryptedPayload.readableBytes() > ClientCnx.getMaxMessageSize()) { cmd.release(); @@ -191,7 +191,7 @@ public OpSendMsg createOpSendMsg() throws IOException { op.setNumMessagesInBatch(numMessagesInBatch); op.setBatchSizeByte(currentBatchSizeBytes); - lowestSequenceId = -1; + lowestSequenceId = -1L; return op; } 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 3d5c58b219075..08101acaebf4a 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 @@ -798,7 +798,7 @@ void ackReceived(ClientCnx cnx, long sequenceId, long ledgerId, long entryId) { } return; } - long expectedSequenceId = op.highestSequenceId > 0 ? op.highestSequenceId : op.sequenceId; + long expectedSequenceId = getHighestSequenceId(op); if (sequenceId > expectedSequenceId) { log.warn("[{}] [{}] Got ack for msg. expecting: {} - got: {} - queue-size: {}", topic, producerName, expectedSequenceId, sequenceId, pendingMessages.size()); @@ -824,7 +824,7 @@ void ackReceived(ClientCnx cnx, long sequenceId, long ledgerId, long entryId) { if (callback) { op = pendingCallbacks.poll(); if (op != null) { - lastSequenceIdPublished = sequenceId; + lastSequenceIdPublished = getHighestSequenceId(op); op.setMessageId(ledgerId, entryId, partitionIndex); try { // Need to protect ourselves from any exception being thrown in the future handler from the @@ -840,6 +840,10 @@ void ackReceived(ClientCnx cnx, long sequenceId, long ledgerId, long entryId) { } } + private long getHighestSequenceId(OpSendMsg op) { + return Math.max(op.highestSequenceId, op.sequenceId); + } + /** * Checks message checksum to retry if message was corrupted while sending to broker. Recomputes checksum of the * message header-payload again. @@ -1430,7 +1434,7 @@ private void processOpSendMsg(OpSendMsg op) { } pendingMessages.put(op); if (op.msg != null) { - lastSequenceIdPushed = op.sequenceId; + lastSequenceIdPushed = getHighestSequenceId(op); } ClientCnx cnx = cnx(); if (isConnected()) { diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java b/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java index dd0c587d9e74b..dcda3df6dbe44 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java @@ -3590,9 +3590,9 @@ public interface MessageMetadataOrBuilder boolean hasTxnidMostBits(); long getTxnidMostBits(); - // optional uint64 last_sequence_id = 24 [default = 0]; - boolean hasLastSequenceId(); - long getLastSequenceId(); + // optional uint64 highest_sequence_id = 24 [default = 0]; + boolean hasHighestSequenceId(); + long getHighestSequenceId(); } public static final class MessageMetadata extends org.apache.pulsar.shaded.com.google.protobuf.v241.GeneratedMessageLite @@ -3953,14 +3953,14 @@ public long getTxnidMostBits() { return txnidMostBits_; } - // optional uint64 last_sequence_id = 24 [default = 0]; - public static final int LAST_SEQUENCE_ID_FIELD_NUMBER = 24; - private long lastSequenceId_; - public boolean hasLastSequenceId() { + // optional uint64 highest_sequence_id = 24 [default = 0]; + public static final int HIGHEST_SEQUENCE_ID_FIELD_NUMBER = 24; + private long highestSequenceId_; + public boolean hasHighestSequenceId() { return ((bitField0_ & 0x00040000) == 0x00040000); } - public long getLastSequenceId() { - return lastSequenceId_; + public long getHighestSequenceId() { + return highestSequenceId_; } private void initFields() { @@ -3985,7 +3985,7 @@ private void initFields() { markerType_ = 0; txnidLeastBits_ = 0L; txnidMostBits_ = 0L; - lastSequenceId_ = 0L; + highestSequenceId_ = 0L; } private byte memoizedIsInitialized = -1; public final boolean isInitialized() { @@ -4092,7 +4092,7 @@ public void writeTo(org.apache.pulsar.common.util.protobuf.ByteBufCodedOutputStr output.writeUInt64(23, txnidMostBits_); } if (((bitField0_ & 0x00040000) == 0x00040000)) { - output.writeUInt64(24, lastSequenceId_); + output.writeUInt64(24, highestSequenceId_); } } @@ -4193,7 +4193,7 @@ public int getSerializedSize() { } if (((bitField0_ & 0x00040000) == 0x00040000)) { size += org.apache.pulsar.shaded.com.google.protobuf.v241.CodedOutputStream - .computeUInt64Size(24, lastSequenceId_); + .computeUInt64Size(24, highestSequenceId_); } memoizedSerializedSize = size; return size; @@ -4350,7 +4350,7 @@ public Builder clear() { bitField0_ = (bitField0_ & ~0x00080000); txnidMostBits_ = 0L; bitField0_ = (bitField0_ & ~0x00100000); - lastSequenceId_ = 0L; + highestSequenceId_ = 0L; bitField0_ = (bitField0_ & ~0x00200000); return this; } @@ -4476,7 +4476,7 @@ public org.apache.pulsar.common.api.proto.PulsarApi.MessageMetadata buildPartial if (((from_bitField0_ & 0x00200000) == 0x00200000)) { to_bitField0_ |= 0x00040000; } - result.lastSequenceId_ = lastSequenceId_; + result.highestSequenceId_ = highestSequenceId_; result.bitField0_ = to_bitField0_; return result; } @@ -4567,8 +4567,8 @@ public Builder mergeFrom(org.apache.pulsar.common.api.proto.PulsarApi.MessageMet if (other.hasTxnidMostBits()) { setTxnidMostBits(other.getTxnidMostBits()); } - if (other.hasLastSequenceId()) { - setLastSequenceId(other.getLastSequenceId()); + if (other.hasHighestSequenceId()) { + setHighestSequenceId(other.getHighestSequenceId()); } return this; } @@ -4736,7 +4736,7 @@ public Builder mergeFrom( } case 192: { bitField0_ |= 0x00200000; - lastSequenceId_ = input.readUInt64(); + highestSequenceId_ = input.readUInt64(); break; } } @@ -5429,23 +5429,23 @@ public Builder clearTxnidMostBits() { return this; } - // optional uint64 last_sequence_id = 24 [default = 0]; - private long lastSequenceId_ ; - public boolean hasLastSequenceId() { + // optional uint64 highest_sequence_id = 24 [default = 0]; + private long highestSequenceId_ ; + public boolean hasHighestSequenceId() { return ((bitField0_ & 0x00200000) == 0x00200000); } - public long getLastSequenceId() { - return lastSequenceId_; + public long getHighestSequenceId() { + return highestSequenceId_; } - public Builder setLastSequenceId(long value) { + public Builder setHighestSequenceId(long value) { bitField0_ |= 0x00200000; - lastSequenceId_ = value; + highestSequenceId_ = value; return this; } - public Builder clearLastSequenceId() { + public Builder clearHighestSequenceId() { bitField0_ = (bitField0_ & ~0x00200000); - lastSequenceId_ = 0L; + highestSequenceId_ = 0L; return this; } @@ -15385,9 +15385,9 @@ public interface CommandSendOrBuilder boolean hasTxnidMostBits(); long getTxnidMostBits(); - // optional uint64 last_sequence_id = 6 [default = 0]; - boolean hasLastSequenceId(); - long getLastSequenceId(); + // optional uint64 highest_sequence_id = 6 [default = 0]; + boolean hasHighestSequenceId(); + long getHighestSequenceId(); } public static final class CommandSend extends org.apache.pulsar.shaded.com.google.protobuf.v241.GeneratedMessageLite @@ -15474,14 +15474,14 @@ public long getTxnidMostBits() { return txnidMostBits_; } - // optional uint64 last_sequence_id = 6 [default = 0]; - public static final int LAST_SEQUENCE_ID_FIELD_NUMBER = 6; - private long lastSequenceId_; - public boolean hasLastSequenceId() { + // optional uint64 highest_sequence_id = 6 [default = 0]; + public static final int HIGHEST_SEQUENCE_ID_FIELD_NUMBER = 6; + private long highestSequenceId_; + public boolean hasHighestSequenceId() { return ((bitField0_ & 0x00000020) == 0x00000020); } - public long getLastSequenceId() { - return lastSequenceId_; + public long getHighestSequenceId() { + return highestSequenceId_; } private void initFields() { @@ -15490,7 +15490,7 @@ private void initFields() { numMessages_ = 1; txnidLeastBits_ = 0L; txnidMostBits_ = 0L; - lastSequenceId_ = 0L; + highestSequenceId_ = 0L; } private byte memoizedIsInitialized = -1; public final boolean isInitialized() { @@ -15533,7 +15533,7 @@ public void writeTo(org.apache.pulsar.common.util.protobuf.ByteBufCodedOutputStr output.writeUInt64(5, txnidMostBits_); } if (((bitField0_ & 0x00000020) == 0x00000020)) { - output.writeUInt64(6, lastSequenceId_); + output.writeUInt64(6, highestSequenceId_); } } @@ -15565,7 +15565,7 @@ public int getSerializedSize() { } if (((bitField0_ & 0x00000020) == 0x00000020)) { size += org.apache.pulsar.shaded.com.google.protobuf.v241.CodedOutputStream - .computeUInt64Size(6, lastSequenceId_); + .computeUInt64Size(6, highestSequenceId_); } memoizedSerializedSize = size; return size; @@ -15690,7 +15690,7 @@ public Builder clear() { bitField0_ = (bitField0_ & ~0x00000008); txnidMostBits_ = 0L; bitField0_ = (bitField0_ & ~0x00000010); - lastSequenceId_ = 0L; + highestSequenceId_ = 0L; bitField0_ = (bitField0_ & ~0x00000020); return this; } @@ -15748,7 +15748,7 @@ public org.apache.pulsar.common.api.proto.PulsarApi.CommandSend buildPartial() { if (((from_bitField0_ & 0x00000020) == 0x00000020)) { to_bitField0_ |= 0x00000020; } - result.lastSequenceId_ = lastSequenceId_; + result.highestSequenceId_ = highestSequenceId_; result.bitField0_ = to_bitField0_; return result; } @@ -15770,8 +15770,8 @@ public Builder mergeFrom(org.apache.pulsar.common.api.proto.PulsarApi.CommandSen if (other.hasTxnidMostBits()) { setTxnidMostBits(other.getTxnidMostBits()); } - if (other.hasLastSequenceId()) { - setLastSequenceId(other.getLastSequenceId()); + if (other.hasHighestSequenceId()) { + setHighestSequenceId(other.getHighestSequenceId()); } return this; } @@ -15837,7 +15837,7 @@ public Builder mergeFrom( } case 48: { bitField0_ |= 0x00000020; - lastSequenceId_ = input.readUInt64(); + highestSequenceId_ = input.readUInt64(); break; } } @@ -15951,23 +15951,23 @@ public Builder clearTxnidMostBits() { return this; } - // optional uint64 last_sequence_id = 6 [default = 0]; - private long lastSequenceId_ ; - public boolean hasLastSequenceId() { + // optional uint64 highest_sequence_id = 6 [default = 0]; + private long highestSequenceId_ ; + public boolean hasHighestSequenceId() { return ((bitField0_ & 0x00000020) == 0x00000020); } - public long getLastSequenceId() { - return lastSequenceId_; + public long getHighestSequenceId() { + return highestSequenceId_; } - public Builder setLastSequenceId(long value) { + public Builder setHighestSequenceId(long value) { bitField0_ |= 0x00000020; - lastSequenceId_ = value; + highestSequenceId_ = value; return this; } - public Builder clearLastSequenceId() { + public Builder clearHighestSequenceId() { bitField0_ = (bitField0_ & ~0x00000020); - lastSequenceId_ = 0L; + highestSequenceId_ = 0L; return this; } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java index 48e70d8f93c0e..4322b3a94ad59 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java @@ -479,8 +479,7 @@ public static ByteBufPair newSend(long producerId, long lowestSequenceId, long h CommandSend.Builder sendBuilder = CommandSend.newBuilder(); sendBuilder.setProducerId(producerId); sendBuilder.setSequenceId(lowestSequenceId); - sendBuilder.setLastSequenceId(highestSequenceId); - sendBuilder.setSequenceId(lowestSequenceId); + sendBuilder.setHighestSequenceId(highestSequenceId); if (numMessages > 1) { sendBuilder.setNumMessages(numMessages); } diff --git a/pulsar-common/src/main/proto/PulsarApi.proto b/pulsar-common/src/main/proto/PulsarApi.proto index 9b99f3539f7d0..f9c88b5f8a09f 100644 --- a/pulsar-common/src/main/proto/PulsarApi.proto +++ b/pulsar-common/src/main/proto/PulsarApi.proto @@ -133,8 +133,8 @@ message MessageMetadata { optional uint64 txnid_least_bits = 22 [default = 0]; optional uint64 txnid_most_bits = 23 [default = 0]; - /// Add last sequence id to support batch message with external sequence id - optional uint64 last_sequence_id = 24 [default = 0]; + /// Add highest sequence id to support batch message with external sequence id + optional uint64 highest_sequence_id = 24 [default = 0]; } message SingleMessageMetadata { @@ -415,8 +415,8 @@ message CommandSend { optional uint64 txnid_least_bits = 4 [default = 0]; optional uint64 txnid_most_bits = 5 [default = 0]; - /// Add last sequence id to support batch message with external sequence id - optional uint64 last_sequence_id = 6 [default = 0]; + /// Add highest sequence id to support batch message with external sequence id + optional uint64 highest_sequence_id = 6 [default = 0]; } message CommandSendReceipt { From ba5e91cc95b05ad9c771cf9c177db59d4585c8c6 Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Tue, 5 Nov 2019 21:14:24 +0800 Subject: [PATCH 11/14] Fix comments --- .../apache/pulsar/broker/service/Producer.java | 6 +++--- .../persistent/MessageDeduplication.java | 18 ++++++++---------- .../pulsar/client/impl/ProducerImpl.java | 4 ++-- 3 files changed, 13 insertions(+), 15 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java index d7366f9083be7..ee9f7eb5e1d4d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java @@ -197,9 +197,9 @@ private void publishMessageToTopic(ByteBuf headersAndPayload, long sequenceId, l System.nanoTime())); } - private void publishMessageToTopic(ByteBuf headersAndPayload, long sequenceId, long lastSequenceId, long batchSize) { + private void publishMessageToTopic(ByteBuf headersAndPayload, long lowestSequenceId, long highestSequenceId, long batchSize) { topic.publishMessage(headersAndPayload, - MessagePublishContext.get(this, sequenceId, lastSequenceId, msgIn, headersAndPayload.readableBytes(), batchSize, + MessagePublishContext.get(this, lowestSequenceId, highestSequenceId, msgIn, headersAndPayload.readableBytes(), batchSize, System.nanoTime())); } @@ -377,7 +377,7 @@ public void run() { // stats rateIn.recordMultipleEvents(batchSize, msgSize); producer.topic.recordAddLatency(System.nanoTime() - startTimeNs, TimeUnit.NANOSECONDS); - long callBackSequenceId = highestSequenceId >= sequenceId ? highestSequenceId : sequenceId;; + long callBackSequenceId = Math.max(highestSequenceId, sequenceId); producer.cnx.ctx().writeAndFlush( Commands.newSendReceipt(producer.producerId, callBackSequenceId, ledgerId, entryId), producer.cnx.ctx().voidPromise()); 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 2794989930c0c..dc4604e8a401c 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 @@ -160,8 +160,7 @@ public void readEntriesComplete(List entries, Object ctx) { MessageMetadata md = Commands.parseMessageMetadata(messageMetadataAndPayload); String producerName = md.getProducerName(); - long sequenceId = md.hasHighestSequenceId() && md.getHighestSequenceId() >= md.getSequenceId() ? - md.getHighestSequenceId() : md.getSequenceId(); + long sequenceId = Math.max(md.getHighestSequenceId(), md.getSequenceId()); highestSequencedPushed.put(producerName, sequenceId); highestSequencedPersisted.put(producerName, sequenceId); @@ -284,8 +283,7 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade String producerName = publishContext.getProducerName(); long sequenceId = publishContext.getSequenceId(); - long lastSequenceId = publishContext.getHighestSequenceId() >= sequenceId ? publishContext.getHighestSequenceId() - : sequenceId; + 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. @@ -293,10 +291,10 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade MessageMetadata md = Commands.parseMessageMetadata(headersAndPayload); producerName = md.getProducerName(); sequenceId = md.getSequenceId(); - lastSequenceId = md.getHighestSequenceId() >= sequenceId ? md.getHighestSequenceId() : md.getSequenceId(); + highestSequenceId = Math.max(md.getHighestSequenceId(), sequenceId); publishContext.setOriginalProducerName(producerName); publishContext.setOriginalSequenceId(sequenceId); - publishContext.setOriginalHighestSequenceId(lastSequenceId); + publishContext.setOriginalHighestSequenceId(highestSequenceId); headersAndPayload.readerIndex(readerIndex); md.recycle(); } @@ -322,7 +320,7 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade return MessageDupStatus.Unknown; } } - highestSequencedPushed.put(producerName, lastSequenceId); + highestSequencedPushed.put(producerName, highestSequenceId); } return MessageDupStatus.NotDup; } @@ -337,15 +335,15 @@ public void recordMessagePersisted(PublishContext publishContext, PositionImpl p String producerName = publishContext.getProducerName(); long sequenceId = publishContext.getSequenceId(); - long lastSequenceId = publishContext.getHighestSequenceId(); + long highestSequenceId = publishContext.getHighestSequenceId(); if (publishContext.getOriginalProducerName() != null) { // In case of replicated messages, this will be different from the current replicator producer name producerName = publishContext.getOriginalProducerName(); sequenceId = publishContext.getOriginalSequenceId(); - lastSequenceId = publishContext.getOriginalHighestSequenceId(); + highestSequenceId = publishContext.getOriginalHighestSequenceId(); } - highestSequencedPersisted.put(producerName, lastSequenceId >= sequenceId ? lastSequenceId : sequenceId); + highestSequencedPersisted.put(producerName, Math.max(highestSequenceId, sequenceId)); if (++snapshotCounter >= snapshotInterval) { snapshotCounter = 0; takeSnapshot(position); 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 08101acaebf4a..f2b48c5df3ed4 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 @@ -393,8 +393,8 @@ public void sendAsync(Message message, SendCallback callback) { if (sequenceId <= lastSequenceIdPublished) { log.warn("Message with sequence id {} is definitely a duplicate", sequenceId); } else { - log.warn("Message with sequence id {} is a definitely a duplicate or not cannot be " + - "determined at this time", sequenceId); + log.info("Message with sequence id {} might be a duplicate but cannot be determined at this time.", + sequenceId); } doBatchSendAndAdd(msg, callback, payload); } else { From b0375ee59c3a19d43d5b4c325d5e2b8db47521c5 Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Tue, 5 Nov 2019 21:16:29 +0800 Subject: [PATCH 12/14] Fix params name --- .../java/org/apache/pulsar/broker/service/Producer.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java index ee9f7eb5e1d4d..a1d8db39e42c9 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java @@ -400,12 +400,12 @@ static MessagePublishContext get(Producer producer, long sequenceId, Rate rateIn return callback; } - static MessagePublishContext get(Producer producer, long sequenceId, long lastSequenceId, Rate rateIn, + static MessagePublishContext get(Producer producer, long lowestSequenceId, long highestSequenceId, Rate rateIn, int msgSize, long batchSize, long startTimeNs) { MessagePublishContext callback = RECYCLER.get(); callback.producer = producer; - callback.sequenceId = sequenceId; - callback.highestSequenceId = lastSequenceId; + callback.sequenceId = lowestSequenceId; + callback.highestSequenceId = highestSequenceId; callback.rateIn = rateIn; callback.msgSize = msgSize; callback.batchSize = batchSize; From cf2f4da0b26313a2f3b86546e8185259b7dbf960 Mon Sep 17 00:00:00 2001 From: lipenghui Date: Wed, 6 Nov 2019 15:06:27 +0800 Subject: [PATCH 13/14] Fix comments --- .../java/org/apache/pulsar/client/impl/ProducerImpl.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) 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 f2b48c5df3ed4..d40805c7af1a1 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 @@ -824,7 +824,7 @@ void ackReceived(ClientCnx cnx, long sequenceId, long ledgerId, long entryId) { if (callback) { op = pendingCallbacks.poll(); if (op != null) { - lastSequenceIdPublished = getHighestSequenceId(op); + lastSequenceIdPublished = Math.max(lastSequenceIdPublished, getHighestSequenceId(op)); op.setMessageId(ledgerId, entryId, partitionIndex); try { // Need to protect ourselves from any exception being thrown in the future handler from the @@ -864,7 +864,7 @@ protected synchronized void recoverChecksumError(ClientCnx cnx, long sequenceId) log.debug("[{}] [{}] Got send failure for timed out msg {}", topic, producerName, sequenceId); } } else { - long expectedSequenceId = op.highestSequenceId > 0 ? op.highestSequenceId : op.sequenceId; + long expectedSequenceId = getHighestSequenceId(op); if (sequenceId == expectedSequenceId) { boolean corrupted = !verifyLocalBufferIsNotCorrupted(op); if (corrupted) { @@ -1434,7 +1434,7 @@ private void processOpSendMsg(OpSendMsg op) { } pendingMessages.put(op); if (op.msg != null) { - lastSequenceIdPushed = getHighestSequenceId(op); + lastSequenceIdPushed = Math.max(lastSequenceIdPushed, getHighestSequenceId(op)); } ClientCnx cnx = cnx(); if (isConnected()) { From d982c8d5f55e408eb8f27df18efcf3be10ffd373 Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Thu, 7 Nov 2019 08:25:35 +0800 Subject: [PATCH 14/14] End with L for long type fields --- .../pulsar/broker/service/Producer.java | 20 +++++++++---------- .../pulsar/client/impl/ProducerImpl.java | 16 +++++++-------- 2 files changed, 18 insertions(+), 18 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java index a1d8db39e42c9..5ead8c944ad77 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java @@ -395,7 +395,7 @@ static MessagePublishContext get(Producer producer, long sequenceId, Rate rateIn callback.msgSize = msgSize; callback.batchSize = batchSize; callback.originalProducerName = null; - callback.originalSequenceId = -1; + callback.originalSequenceId = -1L; callback.startTimeNs = startTimeNs; return callback; } @@ -410,7 +410,7 @@ static MessagePublishContext get(Producer producer, long lowestSequenceId, long callback.msgSize = msgSize; callback.batchSize = batchSize; callback.originalProducerName = null; - callback.originalSequenceId = -1; + callback.originalSequenceId = -1L; callback.startTimeNs = startTimeNs; return callback; } @@ -429,16 +429,16 @@ protected MessagePublishContext newObject(Recycler.Handle public void recycle() { producer = null; - sequenceId = -1; - highestSequenceId = -1; - originalSequenceId = -1; - originalHighestSequenceId = -1; + sequenceId = -1L; + highestSequenceId = -1L; + originalSequenceId = -1L; + originalHighestSequenceId = -1L; rateIn = null; msgSize = 0; - ledgerId = -1; - entryId = -1; - batchSize = 0; - startTimeNs = -1; + ledgerId = -1L; + entryId = -1L; + batchSize = 0L; + startTimeNs = -1L; recyclerHandle.recycle(this); } } 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 d40805c7af1a1..a6365c1ecd27e 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 @@ -143,11 +143,11 @@ public ProducerImpl(PulsarClientImpl client, String topic, ProducerConfiguration long initialSequenceId = conf.getInitialSequenceId(); this.lastSequenceIdPublished = initialSequenceId; this.lastSequenceIdPushed = initialSequenceId; - this.msgIdGenerator = initialSequenceId + 1; + this.msgIdGenerator = initialSequenceId + 1L; } else { - this.lastSequenceIdPublished = -1; - this.lastSequenceIdPushed = -1; - this.msgIdGenerator = 0; + this.lastSequenceIdPublished = -1L; + this.lastSequenceIdPushed = -1L; + this.msgIdGenerator = 0L; } if (conf.isEncryptionEnabled()) { @@ -990,10 +990,10 @@ void recycle() { cmd = null; callback = null; rePopulate = null; - sequenceId = -1; - createdAt = -1; - lowestSequenceId = -1; - highestSequenceId = -1; + sequenceId = -1L; + createdAt = -1L; + lowestSequenceId = -1L; + highestSequenceId = -1L; recyclerHandle.recycle(this); }