From 20b4dc7fedff97ddaf68ce54a07b042506c4c293 Mon Sep 17 00:00:00 2001 From: "wangjiahua.wjh" Date: Thu, 27 Aug 2026 16:39:16 +0800 Subject: [PATCH] [ISSUE #10972] Encode timer message propertiesString after internal properties are cleared --- .../store/timer/TimerMessageStore.java | 4 ++- .../store/timer/TimerMessageStoreTest.java | 25 +++++++++++++++++++ 2 files changed, 28 insertions(+), 1 deletion(-) diff --git a/store/src/main/java/org/apache/rocketmq/store/timer/TimerMessageStore.java b/store/src/main/java/org/apache/rocketmq/store/timer/TimerMessageStore.java index 24d2faa9b61..9608af3213d 100644 --- a/store/src/main/java/org/apache/rocketmq/store/timer/TimerMessageStore.java +++ b/store/src/main/java/org/apache/rocketmq/store/timer/TimerMessageStore.java @@ -1251,7 +1251,6 @@ public MessageExtBrokerInner convertMessage(MessageExt msgExt, boolean needRoll) long tagsCodeValue = MessageExtBrokerInner.tagsString2tagsCode(topicFilterType, msgInner.getTags()); msgInner.setTagsCode(tagsCodeValue); - msgInner.setPropertiesString(MessageDecoder.messageProperties2String(msgExt.getProperties())); msgInner.setSysFlag(msgExt.getSysFlag()); msgInner.setBornTimestamp(msgExt.getBornTimestamp()); @@ -1270,6 +1269,9 @@ public MessageExtBrokerInner convertMessage(MessageExt msgExt, boolean needRoll) MessageAccessor.clearProperty(msgInner, MessageConst.PROPERTY_REAL_TOPIC); MessageAccessor.clearProperty(msgInner, MessageConst.PROPERTY_REAL_QUEUE_ID); } + // Encode after the properties are finalized so that the wire data stays consistent + // with the property map, aligning with TimerMessageRocksDBStore#convertMessage. + msgInner.setPropertiesString(MessageDecoder.messageProperties2String(msgInner.getProperties())); return msgInner; } diff --git a/store/src/test/java/org/apache/rocketmq/store/timer/TimerMessageStoreTest.java b/store/src/test/java/org/apache/rocketmq/store/timer/TimerMessageStoreTest.java index fe1a1177c69..e62075cf67b 100644 --- a/store/src/test/java/org/apache/rocketmq/store/timer/TimerMessageStoreTest.java +++ b/store/src/test/java/org/apache/rocketmq/store/timer/TimerMessageStoreTest.java @@ -179,6 +179,31 @@ private static PutMessageResult transformTimerMessage(TimerMessageStore timerMes return null; } + @Test + public void testConvertMessagePropertiesStringMatchesProperties() throws Exception { + final TimerMessageStore timerMessageStore = createTimerMessageStore(null, true); + + MessageExtBrokerInner msgExt = buildMessage(3000, "TimerTest_testConvertMessage", false); + MessageAccessor.putProperty(msgExt, MessageConst.PROPERTY_REAL_TOPIC, msgExt.getTopic()); + MessageAccessor.putProperty(msgExt, MessageConst.PROPERTY_REAL_QUEUE_ID, "0"); + msgExt.setPropertiesString(MessageDecoder.messageProperties2String(msgExt.getProperties())); + msgExt.setTopic(TimerMessageStore.TIMER_TOPIC); + + // delivered message: internal properties are cleared from both the map and the wire data + MessageExtBrokerInner delivered = timerMessageStore.convertMessage(msgExt, false); + assertEquals("TimerTest_testConvertMessage", delivered.getTopic()); + assertFalse(delivered.getPropertiesString().contains(MessageConst.PROPERTY_REAL_TOPIC)); + assertFalse(delivered.getPropertiesString().contains(MessageConst.PROPERTY_REAL_QUEUE_ID)); + assertEquals(MessageDecoder.messageProperties2String(delivered.getProperties()), delivered.getPropertiesString()); + + // rolled message: keeps REAL_TOPIC and stays consistent between the map and the wire data + MessageExtBrokerInner rolled = timerMessageStore.convertMessage(msgExt, true); + assertEquals(TimerMessageStore.TIMER_TOPIC, rolled.getTopic()); + assertTrue(rolled.getPropertiesString().contains(MessageConst.PROPERTY_REAL_TOPIC)); + assertTrue(rolled.getPropertiesString().contains(MessageConst.PROPERTY_REAL_QUEUE_ID)); + assertEquals(MessageDecoder.messageProperties2String(rolled.getProperties()), rolled.getPropertiesString()); + } + @Test public void testPutTimerMessage() throws Exception { Assume.assumeFalse(MixAll.isWindows());