From 05f2978643950a41f9012f1b27834a037c359e26 Mon Sep 17 00:00:00 2001 From: Ezequiel Lovelle Date: Thu, 27 Dec 2018 20:44:06 -0300 Subject: [PATCH] Support set publish time on broker side Since #3155 there is a real need not to depend on client side clock for publish time field on message metadata. This change sets a default value for publish time field at client side in order to be able to be identified by broker easily and set newly calculated time on it. If time was normally set by the client on publish time, this update will be just ignored by broker. [pulsar-broker] - Read message metadata for produced messages in order to know whether the message was set onto default value and should be updated with time calculated by broker. - Add method to return previous metadata with publish time set by broker. [pulsar-client] - ProducerImpl client side set a default fixed size value for publish time. [pulsar-common] - Add modifyMessageMetadata method on Commands in order to update new modified fields on payload. The default constant value used to define the undefined number of milliseconds since epoch is going to be 13 decimal digits, the method currentTimeMillis will always return 13 decimal digits assuming the broker has the current time set correctly. So we are safe with the assumption that the return value is 13 decimal digits for more than the next two centuries. --- .../pulsar/broker/service/Producer.java | 14 ++++++-- .../pulsar/client/impl/ProducerImpl.java | 3 +- .../apache/pulsar/common/api/Commands.java | 34 +++++++++++++++++++ .../pulsar/common/naming/Constants.java | 1 + 4 files changed, 49 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 eb087bca136b4..a83278a009a36 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 @@ -43,6 +43,7 @@ import org.apache.pulsar.common.api.Commands; import org.apache.pulsar.common.api.proto.PulsarApi.MessageMetadata; import org.apache.pulsar.common.api.proto.PulsarApi.ServerError; +import org.apache.pulsar.common.naming.Constants; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.NonPersistentPublisherStats; import org.apache.pulsar.common.policies.data.PublisherStats; @@ -148,12 +149,17 @@ public void publishMessage(long producerId, long sequenceId, ByteBuf headersAndP return; } - if (topic.isEncryptionRequired()) { + headersAndPayload.markReaderIndex(); + MessageMetadata msgMetadata = Commands.parseMessageMetadata(headersAndPayload); + headersAndPayload.resetReaderIndex(); + if (msgMetadata.getPublishTime() == Constants.PUBLISH_TIME_UNSET_MS) { headersAndPayload.markReaderIndex(); - MessageMetadata msgMetadata = Commands.parseMessageMetadata(headersAndPayload); + Commands.modifyMessageMetadata(headersAndPayload, newPublishOnMetadata(msgMetadata)); headersAndPayload.resetReaderIndex(); + } + if (topic.isEncryptionRequired()) { // Check whether the message is encrypted or not if (msgMetadata.getEncryptionKeysCount() < 1) { log.warn("[{}] Messages must be encrypted", getTopic().getName()); @@ -171,6 +177,10 @@ public void publishMessage(long producerId, long sequenceId, ByteBuf headersAndP MessagePublishContext.get(this, sequenceId, msgIn, headersAndPayload.readableBytes(), batchSize)); } + private static MessageMetadata newPublishOnMetadata(MessageMetadata msgMetadata) { + return msgMetadata.toBuilder().setPublishTime(System.currentTimeMillis()).build(); + } + private boolean verifyChecksum(ByteBuf headersAndPayload) { if (hasChecksum(headersAndPayload)) { int readerIndex = headersAndPayload.readerIndex(); 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 88f9b1c4196da..ec426163b6d4f 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 @@ -67,6 +67,7 @@ import org.apache.pulsar.common.api.proto.PulsarApi.ProtocolVersion; import org.apache.pulsar.common.compression.CompressionCodec; import org.apache.pulsar.common.compression.CompressionCodecProvider; +import org.apache.pulsar.common.naming.Constants; import org.apache.pulsar.common.schema.SchemaInfo; import org.apache.pulsar.common.schema.SchemaType; import org.apache.pulsar.common.util.DateFormatter; @@ -344,7 +345,7 @@ public void sendAsync(Message message, SendCallback callback) { sequenceId = msgMetadataBuilder.getSequenceId(); } if (!msgMetadataBuilder.hasPublishTime()) { - msgMetadataBuilder.setPublishTime(System.currentTimeMillis()); + msgMetadataBuilder.setPublishTime(Constants.PUBLISH_TIME_UNSET_MS); checkArgument(!msgMetadataBuilder.hasProducerName()); diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/api/Commands.java b/pulsar-common/src/main/java/org/apache/pulsar/common/api/Commands.java index 2adb274697289..665a52a60436b 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/api/Commands.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/api/Commands.java @@ -262,6 +262,40 @@ public static void skipChecksumIfPresent(ByteBuf buffer) { } } + public static void modifyMessageMetadata(ByteBuf buffer, MessageMetadata newMsgMetadata) { + try { + int checksumReaderIndex = buffer.readerIndex(); + skipChecksumIfPresent(buffer); + int metadataSize = (int) buffer.readUnsignedInt(); + int newMetadataSize = newMsgMetadata.getSerializedSize(); + + if (metadataSize != newMetadataSize) { + throw new RuntimeException("Write new metadata with different size from previous is not yet supported"); + } + + buffer.markWriterIndex(); + buffer.writerIndex(buffer.readerIndex()); + ByteBufCodedOutputStream outputStream = ByteBufCodedOutputStream.get(buffer); + newMsgMetadata.writeTo(outputStream); + buffer.resetWriterIndex(); + + buffer.readerIndex(checksumReaderIndex); + if (hasChecksum(buffer)) { + int currentChecksum = readChecksum(buffer); + int computedChecksum = computeChecksum(buffer); + // set new computed checksum + if (currentChecksum != computedChecksum) { + buffer.setInt(checksumReaderIndex+2, computedChecksum); + } + } + + outputStream.recycle(); + newMsgMetadata.recycle(); + } catch (IOException e) { + throw new RuntimeException(e); + } + } + public static MessageMetadata parseMessageMetadata(ByteBuf buffer) { try { // initially reader-index may point to start_of_checksum : increment reader-index to start_of_metadata to parse diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/naming/Constants.java b/pulsar-common/src/main/java/org/apache/pulsar/common/naming/Constants.java index aed20fa556701..28a82a7b0ddc6 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/naming/Constants.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/naming/Constants.java @@ -20,6 +20,7 @@ public class Constants { + public static final long PUBLISH_TIME_UNSET_MS = 1000000000000L; public static final String GLOBAL_CLUSTER = "global"; private Constants() {}