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 224e093baf112..59bd3b4b073ce 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 @@ -472,7 +472,7 @@ public static void skipMessageMetadata(ByteBuf buffer) { } public static long getEntryTimestamp(ByteBuf headersAndPayloadWithBrokerEntryMetadata) throws IOException { - // get broker timestamp first if BrokerEntryMetadata is enabled with AppendBrokerTimestampMetadataInterceptor + // get broker timestamp first if exists BrokerEntryMetadata brokerEntryMetadata = Commands.parseBrokerEntryMetadataIfExist(headersAndPayloadWithBrokerEntryMetadata); if (brokerEntryMetadata != null && brokerEntryMetadata.hasBrokerTimestamp()) { @@ -1716,6 +1716,9 @@ public static ByteBuf addBrokerEntryMetadata(ByteBuf headerAndPayload, interceptor.interceptWithNumberOfMessages(brokerEntryMetadata, numberOfMessages); } } + if (!brokerEntryMetadata.hasBrokerTimestamp()) { + brokerEntryMetadata.setBrokerTimestamp(System.currentTimeMillis()); + } int brokerMetaSize = brokerEntryMetadata.getSerializedSize(); ByteBuf brokerMeta = diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/protocol/CommandUtilsTests.java b/pulsar-common/src/test/java/org/apache/pulsar/common/protocol/CommandUtilsTests.java index 4523a5cc97b37..f508eb3d216eb 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/protocol/CommandUtilsTests.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/protocol/CommandUtilsTests.java @@ -20,6 +20,8 @@ import static org.apache.pulsar.common.protocol.Commands.serializeMetadataAndPayload; import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertNotEquals; +import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertTrue; import io.netty.buffer.ByteBuf; @@ -151,6 +153,26 @@ public void testAddBrokerEntryMetadata() throws Exception { assertTrue(new String(content, StandardCharsets.UTF_8).endsWith(data)); } + @Test + public void testAlwaysSetBrokerTimestamp() { + int MOCK_BATCH_SIZE = 10; + String data = "test-message"; + ByteBuf byteBuf = PulsarByteBufAllocator.DEFAULT.buffer(data.length(), data.length()); + byteBuf.writeBytes(data.getBytes(StandardCharsets.UTF_8)); + ByteBuf dataWithBrokerEntryMetadata = + Commands.addBrokerEntryMetadata(byteBuf, new HashSet<>(), MOCK_BATCH_SIZE); + BrokerEntryMetadata entryMetadata = Commands.parseBrokerEntryMetadataIfExist(dataWithBrokerEntryMetadata); + assertNotNull(entryMetadata); + assertTrue(entryMetadata.hasBrokerTimestamp()); + long timestamp = entryMetadata.getBrokerTimestamp(); + // add again, the timestamp should change + dataWithBrokerEntryMetadata = + Commands.addBrokerEntryMetadata(byteBuf, new HashSet<>(), MOCK_BATCH_SIZE); + entryMetadata = Commands.parseBrokerEntryMetadataIfExist(dataWithBrokerEntryMetadata); + assertNotNull(entryMetadata); + assertNotEquals(timestamp, entryMetadata.getBrokerTimestamp()); + } + @Test public void testSkipBrokerEntryMetadata() throws Exception { String data = "test-message";