From 501c3a14833b6aa2f64e0679f89c94858a0e6f13 Mon Sep 17 00:00:00 2001 From: ccwss <1782935682@qq.com> Date: Fri, 5 Sep 2025 10:55:08 +0800 Subject: [PATCH 1/2] fix error when update topic --- .../main/java/org/apache/rocketmq/common/TopicConfig.java | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/common/src/main/java/org/apache/rocketmq/common/TopicConfig.java b/common/src/main/java/org/apache/rocketmq/common/TopicConfig.java index 0bf64905a03..cf72aa13b0d 100644 --- a/common/src/main/java/org/apache/rocketmq/common/TopicConfig.java +++ b/common/src/main/java/org/apache/rocketmq/common/TopicConfig.java @@ -22,6 +22,7 @@ import java.util.HashMap; import java.util.Map; import java.util.Objects; +import org.apache.rocketmq.common.attribute.AttributeParser; import org.apache.rocketmq.common.attribute.TopicMessageType; import org.apache.rocketmq.common.constant.PermName; @@ -203,7 +204,9 @@ public TopicMessageType getTopicMessageType() { if (attributes == null) { return TopicMessageType.NORMAL; } - String content = attributes.get(TOPIC_MESSAGE_TYPE_ATTRIBUTE.getName()); + String content = attributes.get(TOPIC_MESSAGE_TYPE_ATTRIBUTE.getName()) == null + ? attributes.get(AttributeParser.ATTR_ADD_PLUS_SIGN + TOPIC_MESSAGE_TYPE_ATTRIBUTE.getName()) + : attributes.get(TOPIC_MESSAGE_TYPE_ATTRIBUTE.getName()); if (content == null) { return TopicMessageType.NORMAL; } @@ -212,7 +215,7 @@ public TopicMessageType getTopicMessageType() { @JSONField(serialize = false, deserialize = false) public void setTopicMessageType(TopicMessageType topicMessageType) { - attributes.put(TOPIC_MESSAGE_TYPE_ATTRIBUTE.getName(), topicMessageType.getValue()); + attributes.put(AttributeParser.ATTR_ADD_PLUS_SIGN + TOPIC_MESSAGE_TYPE_ATTRIBUTE.getName(), topicMessageType.getValue()); } @Override From 95f4d4197c6faf0f9f6a953194e098925f09bd56 Mon Sep 17 00:00:00 2001 From: ccwss <1782935682@qq.com> Date: Fri, 5 Sep 2025 11:05:17 +0800 Subject: [PATCH 2/2] fix test --- .../test/java/org/apache/rocketmq/common/TopicConfigTest.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/common/src/test/java/org/apache/rocketmq/common/TopicConfigTest.java b/common/src/test/java/org/apache/rocketmq/common/TopicConfigTest.java index 3df93a0bfb3..d2b5f958fbc 100644 --- a/common/src/test/java/org/apache/rocketmq/common/TopicConfigTest.java +++ b/common/src/test/java/org/apache/rocketmq/common/TopicConfigTest.java @@ -40,12 +40,12 @@ public void testEncode() { topicConfig.setTopicMessageType(TopicMessageType.FIFO); String encode = topicConfig.encode(); - assertThat(encode).isEqualTo("topic 8 8 6 SINGLE_TAG {\"message.type\":\"FIFO\"}"); + assertThat(encode).isEqualTo("topic 8 8 6 SINGLE_TAG {\"+message.type\":\"FIFO\"}"); } @Test public void testDecode() { - String encode = "topic 8 8 6 SINGLE_TAG {\"message.type\":\"FIFO\"}"; + String encode = "topic 8 8 6 SINGLE_TAG {\"+message.type\":\"FIFO\"}"; TopicConfig decodeTopicConfig = new TopicConfig(); decodeTopicConfig.decode(encode);