This is an automated email from the ASF dual-hosted git repository.

lwclover pushed a commit to branch develop
in repository https://gitbox.apache.org/repos/asf/rocketmq.git


The following commit(s) were added to refs/heads/develop by this push:
     new 4dad755fdc [ISSUE #10980] Encode and decode POP checkpoint and ack 
messages without the intermediate JSON string (#10981)
4dad755fdc is described below

commit 4dad755fdc07a81650c04c43c211952d1a27b2a0
Author: Jiahua Wang <[email protected]>
AuthorDate: Tue Sep 15 08:42:36 2026 +0800

    [ISSUE #10980] Encode and decode POP checkpoint and ack messages without 
the intermediate JSON string (#10981)
    
    Co-authored-by: wangjiahua.wjh <[email protected]>
---
 .../broker/processor/AckMessageProcessor.java      |  3 +-
 .../processor/ChangeInvisibleTimeProcessor.java    |  5 +-
 .../broker/processor/PopBufferMergeService.java    |  4 +-
 .../broker/processor/PopMessageProcessor.java      |  3 +-
 .../broker/processor/PopReviveService.java         | 21 ++++-----
 .../org/apache/rocketmq/store/pop/AckMsgTest.java  | 21 +++++++++
 .../apache/rocketmq/store/pop/BatchAckMsgTest.java | 27 ++++++++++-
 .../rocketmq/store/pop/PopCheckPointTest.java      | 53 ++++++++++++++++++++++
 8 files changed, 114 insertions(+), 23 deletions(-)

diff --git 
a/broker/src/main/java/org/apache/rocketmq/broker/processor/AckMessageProcessor.java
 
b/broker/src/main/java/org/apache/rocketmq/broker/processor/AckMessageProcessor.java
index 65f5f79aec..036a7040e4 100644
--- 
a/broker/src/main/java/org/apache/rocketmq/broker/processor/AckMessageProcessor.java
+++ 
b/broker/src/main/java/org/apache/rocketmq/broker/processor/AckMessageProcessor.java
@@ -19,7 +19,6 @@ package org.apache.rocketmq.broker.processor;
 import com.alibaba.fastjson2.JSON;
 import io.netty.channel.Channel;
 import io.netty.channel.ChannelHandlerContext;
-import java.nio.charset.StandardCharsets;
 import java.util.BitSet;
 import org.apache.commons.lang3.StringUtils;
 import org.apache.rocketmq.broker.BrokerController;
@@ -285,7 +284,7 @@ public class AckMessageProcessor implements 
NettyRequestProcessor {
 
         MessageExtBrokerInner msgInner = new MessageExtBrokerInner();
         msgInner.setTopic(reviveTopic);
-        
msgInner.setBody(JSON.toJSONString(ackMsg).getBytes(StandardCharsets.UTF_8));
+        msgInner.setBody(JSON.toJSONBytes(ackMsg));
         msgInner.setQueueId(rqId);
         if (ackMsg instanceof BatchAckMsg) {
             msgInner.setTags(PopAckConstants.BATCH_ACK_TAG);
diff --git 
a/broker/src/main/java/org/apache/rocketmq/broker/processor/ChangeInvisibleTimeProcessor.java
 
b/broker/src/main/java/org/apache/rocketmq/broker/processor/ChangeInvisibleTimeProcessor.java
index 02deeb18a7..fce4a74bb6 100644
--- 
a/broker/src/main/java/org/apache/rocketmq/broker/processor/ChangeInvisibleTimeProcessor.java
+++ 
b/broker/src/main/java/org/apache/rocketmq/broker/processor/ChangeInvisibleTimeProcessor.java
@@ -19,7 +19,6 @@ package org.apache.rocketmq.broker.processor;
 import com.alibaba.fastjson2.JSON;
 import io.netty.channel.Channel;
 import io.netty.channel.ChannelHandlerContext;
-import java.nio.charset.StandardCharsets;
 import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.TimeUnit;
 import org.apache.commons.lang3.StringUtils;
@@ -278,7 +277,7 @@ public class ChangeInvisibleTimeProcessor implements 
NettyRequestProcessor {
         }
 
         msgInner.setTopic(reviveTopic);
-        
msgInner.setBody(JSON.toJSONString(ackMsg).getBytes(StandardCharsets.UTF_8));
+        msgInner.setBody(JSON.toJSONBytes(ackMsg));
         msgInner.setQueueId(rqId);
         msgInner.setTags(PopAckConstants.ACK_TAG);
         msgInner.setBornTimestamp(System.currentTimeMillis());
@@ -322,7 +321,7 @@ public class ChangeInvisibleTimeProcessor implements 
NettyRequestProcessor {
         ck.setBrokerName(ExtraInfoUtil.getBrokerName(extraInfo));
         ck.setSuspend(requestHeader.isSuspend());
 
-        
msgInner.setBody(JSON.toJSONString(ck).getBytes(StandardCharsets.UTF_8));
+        msgInner.setBody(JSON.toJSONBytes(ck));
         msgInner.setQueueId(reviveQid);
         msgInner.setTags(PopAckConstants.CK_TAG);
         msgInner.setBornTimestamp(System.currentTimeMillis());
diff --git 
a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopBufferMergeService.java
 
b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopBufferMergeService.java
index 5373eaea33..13869e958e 100644
--- 
a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopBufferMergeService.java
+++ 
b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopBufferMergeService.java
@@ -668,7 +668,7 @@ public class PopBufferMergeService extends ServiceThread {
         ackMsg.setPopTime(point.getPopTime());
         ackMsg.setBrokerName(point.getBrokerName());
         msgInner.setTopic(popMessageProcessor.getReviveTopic());
-        
msgInner.setBody(JSON.toJSONString(ackMsg).getBytes(DataConverter.CHARSET_UTF8));
+        msgInner.setBody(JSON.toJSONBytes(ackMsg));
         msgInner.setQueueId(pointWrapper.getReviveQueueId());
         msgInner.setTags(PopAckConstants.ACK_TAG);
         msgInner.setBornTimestamp(System.currentTimeMillis());
@@ -724,7 +724,7 @@ public class PopBufferMergeService extends ServiceThread {
         batchAckMsg.setQueueId(point.getQueueId());
         batchAckMsg.setPopTime(point.getPopTime());
         msgInner.setTopic(popMessageProcessor.getReviveTopic());
-        
msgInner.setBody(JSON.toJSONString(batchAckMsg).getBytes(DataConverter.CHARSET_UTF8));
+        msgInner.setBody(JSON.toJSONBytes(batchAckMsg));
         msgInner.setQueueId(pointWrapper.getReviveQueueId());
         msgInner.setTags(PopAckConstants.BATCH_ACK_TAG);
         msgInner.setBornTimestamp(System.currentTimeMillis());
diff --git 
a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopMessageProcessor.java
 
b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopMessageProcessor.java
index 55cabe6f5e..779e6fef2d 100644
--- 
a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopMessageProcessor.java
+++ 
b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopMessageProcessor.java
@@ -78,7 +78,6 @@ import org.apache.rocketmq.store.pop.BatchAckMsg;
 import org.apache.rocketmq.store.pop.PopCheckPoint;
 
 import java.nio.ByteBuffer;
-import java.nio.charset.StandardCharsets;
 import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
@@ -937,7 +936,7 @@ public class PopMessageProcessor implements 
NettyRequestProcessor {
         MessageExtBrokerInner msgInner = new MessageExtBrokerInner();
 
         msgInner.setTopic(reviveTopic);
-        
msgInner.setBody(JSON.toJSONString(ck).getBytes(StandardCharsets.UTF_8));
+        msgInner.setBody(JSON.toJSONBytes(ck));
         msgInner.setQueueId(reviveQid);
         msgInner.setTags(PopAckConstants.CK_TAG);
         msgInner.setBornTimestamp(System.currentTimeMillis());
diff --git 
a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java
 
b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java
index 07f16e9896..67abcfa907 100644
--- 
a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java
+++ 
b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java
@@ -385,11 +385,10 @@ public class PopReviveService extends ServiceThread {
             }
             for (MessageExt messageExt : messageExts) {
                 if (PopAckConstants.CK_TAG.equals(messageExt.getTags())) {
-                    String raw = new String(messageExt.getBody(), 
DataConverter.CHARSET_UTF8);
-                    if (brokerController.getBrokerConfig().isEnablePopLog()) {
-                        POP_LOGGER.info("reviveQueueId={},find ck, offset:{}, 
raw : {}", messageExt.getQueueId(), messageExt.getQueueOffset(), raw);
+                    if (brokerController.getBrokerConfig().isEnablePopLog() && 
POP_LOGGER.isInfoEnabled()) {
+                        POP_LOGGER.info("reviveQueueId={},find ck, offset:{}, 
raw : {}", messageExt.getQueueId(), messageExt.getQueueOffset(), new 
String(messageExt.getBody(), DataConverter.CHARSET_UTF8));
                     }
-                    PopCheckPoint point = JSON.parseObject(raw, 
PopCheckPoint.class);
+                    PopCheckPoint point = 
JSON.parseObject(messageExt.getBody(), PopCheckPoint.class);
                     if (point.getTopic() == null || point.getCId() == null) {
                         continue;
                     }
@@ -400,11 +399,10 @@ public class PopReviveService extends ServiceThread {
                         firstRt = point.getReviveTime();
                     }
                 } else if 
(PopAckConstants.ACK_TAG.equals(messageExt.getTags())) {
-                    String raw = new String(messageExt.getBody(), 
StandardCharsets.UTF_8);
-                    if (brokerController.getBrokerConfig().isEnablePopLog()) {
-                        POP_LOGGER.info("reviveQueueId={}, find ack, 
offset:{}, raw : {}", messageExt.getQueueId(), messageExt.getQueueOffset(), 
raw);
+                    if (brokerController.getBrokerConfig().isEnablePopLog() && 
POP_LOGGER.isInfoEnabled()) {
+                        POP_LOGGER.info("reviveQueueId={}, find ack, 
offset:{}, raw : {}", messageExt.getQueueId(), messageExt.getQueueOffset(), new 
String(messageExt.getBody(), StandardCharsets.UTF_8));
                     }
-                    AckMsg ackMsg = JSON.parseObject(raw, AckMsg.class);
+                    AckMsg ackMsg = JSON.parseObject(messageExt.getBody(), 
AckMsg.class);
                     
brokerController.getBrokerMetricsManager().getPopMetricsManager().incPopReviveAckGetCount(ackMsg,
 queueId);
                     String brokerName = 
StringUtils.isNotBlank(ackMsg.getBrokerName()) ?
                         ackMsg.getBrokerName() : 
brokerController.getBrokerConfig().getBrokerName();
@@ -426,12 +424,11 @@ public class PopReviveService extends ServiceThread {
                         }
                     }
                 } else if 
(PopAckConstants.BATCH_ACK_TAG.equals(messageExt.getTags())) {
-                    String raw = new String(messageExt.getBody(), 
StandardCharsets.UTF_8);
-                    if (brokerController.getBrokerConfig().isEnablePopLog()) {
-                        POP_LOGGER.info("reviveQueueId={}, find batch ack, 
offset:{}, raw : {}", messageExt.getQueueId(), messageExt.getQueueOffset(), 
raw);
+                    if (brokerController.getBrokerConfig().isEnablePopLog() && 
POP_LOGGER.isInfoEnabled()) {
+                        POP_LOGGER.info("reviveQueueId={}, find batch ack, 
offset:{}, raw : {}", messageExt.getQueueId(), messageExt.getQueueOffset(), new 
String(messageExt.getBody(), StandardCharsets.UTF_8));
                     }
 
-                    BatchAckMsg bAckMsg = JSON.parseObject(raw, 
BatchAckMsg.class);
+                    BatchAckMsg bAckMsg = 
JSON.parseObject(messageExt.getBody(), BatchAckMsg.class);
                     
brokerController.getBrokerMetricsManager().getPopMetricsManager().incPopReviveAckGetCount(bAckMsg,
 queueId);
                     String brokerName = 
StringUtils.isNotBlank(bAckMsg.getBrokerName()) ?
                         bAckMsg.getBrokerName() : 
brokerController.getBrokerConfig().getBrokerName();
diff --git a/store/src/test/java/org/apache/rocketmq/store/pop/AckMsgTest.java 
b/store/src/test/java/org/apache/rocketmq/store/pop/AckMsgTest.java
index 13df028e83..14a147ea62 100644
--- a/store/src/test/java/org/apache/rocketmq/store/pop/AckMsgTest.java
+++ b/store/src/test/java/org/apache/rocketmq/store/pop/AckMsgTest.java
@@ -18,6 +18,7 @@
 package org.apache.rocketmq.store.pop;
 
 import com.alibaba.fastjson2.JSON;
+import java.nio.charset.StandardCharsets;
 import org.junit.Assert;
 import org.junit.Test;
 
@@ -48,4 +49,24 @@ public class AckMsgTest {
         Assert.assertEquals(ackMsg1.getAckOffset(), ackMsg2.getAckOffset());
         Assert.assertEquals(ackMsg1.getPopTime(), ackMsg2.getPopTime());
     }
+
+    @Test
+    public void testToJsonBytesMatchesJsonStringBytes() {
+        AckMsg ackMsg = new AckMsg();
+        ackMsg.setBrokerName("broker-a");
+        ackMsg.setTopic("topic-\u4e2d\u6587");
+        ackMsg.setConsumerGroup("group");
+        ackMsg.setQueueId(3);
+        ackMsg.setStartOffset(200L);
+        ackMsg.setAckOffset(100L);
+        ackMsg.setPopTime(1670212915531L);
+
+        byte[] direct = JSON.toJSONBytes(ackMsg);
+        
Assert.assertArrayEquals(JSON.toJSONString(ackMsg).getBytes(StandardCharsets.UTF_8),
 direct);
+
+        AckMsg decoded = JSON.parseObject(direct, AckMsg.class);
+        Assert.assertEquals(ackMsg.getTopic(), decoded.getTopic());
+        Assert.assertEquals(ackMsg.getAckOffset(), decoded.getAckOffset());
+        Assert.assertEquals(ackMsg.getPopTime(), decoded.getPopTime());
+    }
 }
\ No newline at end of file
diff --git 
a/store/src/test/java/org/apache/rocketmq/store/pop/BatchAckMsgTest.java 
b/store/src/test/java/org/apache/rocketmq/store/pop/BatchAckMsgTest.java
index 0a1bc714cf..53b16b4dd1 100644
--- a/store/src/test/java/org/apache/rocketmq/store/pop/BatchAckMsgTest.java
+++ b/store/src/test/java/org/apache/rocketmq/store/pop/BatchAckMsgTest.java
@@ -18,12 +18,14 @@
 package org.apache.rocketmq.store.pop;
 
 import com.alibaba.fastjson2.JSON;
-import org.junit.Assert;
-import org.junit.Test;
 
+import java.nio.charset.StandardCharsets;
 import java.util.ArrayList;
 import java.util.List;
 
+import org.junit.Assert;
+import org.junit.Test;
+
 public class BatchAckMsgTest {
 
     @Test
@@ -54,4 +56,25 @@ public class BatchAckMsgTest {
         Assert.assertEquals(batchAckMsg1.getStartOffset(), 
batchAckMsg2.getStartOffset());
         Assert.assertEquals(batchAckMsg1.getPopTime(), 
batchAckMsg2.getPopTime());
     }
+
+    @Test
+    public void testToJsonBytesMatchesJsonStringBytes() {
+        BatchAckMsg batchAckMsg = new BatchAckMsg();
+        List<Long> aol = new ArrayList<>(2);
+        aol.add(100L);
+        aol.add(101L);
+        batchAckMsg.setAckOffsetList(aol);
+        batchAckMsg.setStartOffset(200L);
+        batchAckMsg.setConsumerGroup("group");
+        batchAckMsg.setTopic("topic-\u4e2d\u6587");
+        batchAckMsg.setQueueId(3);
+        batchAckMsg.setPopTime(1679454922000L);
+
+        byte[] direct = JSON.toJSONBytes(batchAckMsg);
+        
Assert.assertArrayEquals(JSON.toJSONString(batchAckMsg).getBytes(StandardCharsets.UTF_8),
 direct);
+
+        BatchAckMsg decoded = JSON.parseObject(direct, BatchAckMsg.class);
+        Assert.assertEquals(batchAckMsg.getAckOffsetList(), 
decoded.getAckOffsetList());
+        Assert.assertEquals(batchAckMsg.getTopic(), decoded.getTopic());
+    }
 }
diff --git 
a/store/src/test/java/org/apache/rocketmq/store/pop/PopCheckPointTest.java 
b/store/src/test/java/org/apache/rocketmq/store/pop/PopCheckPointTest.java
new file mode 100644
index 0000000000..4b74a42b67
--- /dev/null
+++ b/store/src/test/java/org/apache/rocketmq/store/pop/PopCheckPointTest.java
@@ -0,0 +1,53 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.store.pop;
+
+import com.alibaba.fastjson2.JSON;
+import java.nio.charset.StandardCharsets;
+import org.junit.Assert;
+import org.junit.Test;
+
+public class PopCheckPointTest {
+
+    @Test
+    public void testToJsonBytesMatchesJsonStringBytes() {
+        PopCheckPoint ck = new PopCheckPoint();
+        ck.setTopic("topic-\u4e2d\u6587");
+        ck.setCId("group");
+        ck.setQueueId(3);
+        ck.setStartOffset(200L);
+        ck.setPopTime(1670212915531L);
+        ck.setInvisibleTime(60000L);
+        ck.setBitMap(5);
+        ck.setNum((byte) 2);
+        ck.setBrokerName("broker-a");
+        ck.addDiff(1);
+        ck.addDiff(3);
+        ck.setRePutTimes("1");
+
+        byte[] direct = JSON.toJSONBytes(ck);
+        
Assert.assertArrayEquals(JSON.toJSONString(ck).getBytes(StandardCharsets.UTF_8),
 direct);
+
+        PopCheckPoint decoded = JSON.parseObject(direct, PopCheckPoint.class);
+        Assert.assertEquals(ck.getTopic(), decoded.getTopic());
+        Assert.assertEquals(ck.getCId(), decoded.getCId());
+        Assert.assertEquals(ck.getStartOffset(), decoded.getStartOffset());
+        Assert.assertEquals(ck.getPopTime(), decoded.getPopTime());
+        Assert.assertEquals(ck.getQueueOffsetDiff(), 
decoded.getQueueOffsetDiff());
+        Assert.assertEquals(ck.getBitMap(), decoded.getBitMap());
+    }
+}

Reply via email to