This is an automated email from the ASF dual-hosted git repository.
fuyou001 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 b37e2bbacd [ISSUE #10639] Reuse a per-thread scratch buffer in
CommitLog.checkMessageAndReturnSize (#10640)
b37e2bbacd is described below
commit b37e2bbacd2c899112f22028b139dfdfb10da712
Author: Jiahua Wang <[email protected]>
AuthorDate: Thu Jul 23 23:55:33 2026 +0800
[ISSUE #10639] Reuse a per-thread scratch buffer in
CommitLog.checkMessageAndReturnSize (#10640)
* [ISSUE #10639] Reuse a per-thread scratch buffer in
CommitLog.checkMessageAndReturnSize
* [ISSUE #10639] Reject corrupt negative totalSize before allocation in
checkMessageAndReturnSize
* [ISSUE #10639] Add configurable reuse-buffer cap and unit tests for
checkMessageAndReturnSize
* [ISSUE #10639] Default the check-message reuse-buffer cap to 1 MiB
---------
Co-authored-by: wangjiahua.wjh <[email protected]>
---
.../java/org/apache/rocketmq/store/CommitLog.java | 34 +++++-
.../rocketmq/store/config/MessageStoreConfig.java | 17 +++
.../store/CheckMessageBufferReuseTest.java | 129 +++++++++++++++++++++
3 files changed, 178 insertions(+), 2 deletions(-)
diff --git a/store/src/main/java/org/apache/rocketmq/store/CommitLog.java
b/store/src/main/java/org/apache/rocketmq/store/CommitLog.java
index d2f2da8b7d..260f021aed 100644
--- a/store/src/main/java/org/apache/rocketmq/store/CommitLog.java
+++ b/store/src/main/java/org/apache/rocketmq/store/CommitLog.java
@@ -109,6 +109,14 @@ public class CommitLog implements Swappable {
private final boolean enabledAppendPropCRC;
+ // Per-thread reusable scratch buffer for checkMessageAndReturnSize. That
method reads
+ // body/topic/properties into this buffer transiently (contents are copied
out or CRC'd
+ // before the method returns), so per-thread reuse is safe and avoids
allocating a fresh
+ // message-sized byte[] on every dispatched message. Buffers are grow-only
up to
+ // maxCheckMessageReuseBufferSize; larger sizes (oversized messages or a
corrupt length)
+ // fall back to a transient array to avoid pinning a huge buffer
permanently.
+ private final ThreadLocal<byte[]> checkMessageBuffer =
ThreadLocal.withInitial(() -> new byte[0]);
+
public CommitLog(final DefaultMessageStore messageStore) {
String storePath =
messageStore.getMessageStoreConfig().getStorePathCommitLog();
RunningFlags runningFlags =
messageStore.getMessageStoreConfig().isEnableRunningFlagsInFlush()
@@ -432,6 +440,28 @@ public class CommitLog implements Swappable {
}
}
+ /**
+ * Returns a scratch buffer of at least {@code totalSize} bytes for {@link
#checkMessageAndReturnSize}.
+ * The buffer is used purely transiently there (each field is read into it
and copied out / CRC'd
+ * before the next use), so a per-thread reusable buffer is safe and
removes a per-message byte[]
+ * allocation. Sizes within {@code maxCheckMessageReuseBufferSize} are
reused grow-only; larger sizes
+ * return a transient array so an oversized buffer is never pinned. A
negative (corrupt) totalSize
+ * is rejected by the caller before reaching here.
+ */
+ // Package-private for testing (CheckMessageBufferReuseTest).
+ byte[] borrowCheckMessageBuffer(final int totalSize) {
+ int reuseCap =
this.defaultMessageStore.getMessageStoreConfig().getMaxCheckMessageReuseBufferSize();
+ if (totalSize > reuseCap) {
+ return new byte[totalSize];
+ }
+ byte[] buffer = this.checkMessageBuffer.get();
+ if (buffer.length < totalSize) {
+ buffer = new byte[totalSize];
+ this.checkMessageBuffer.set(buffer);
+ }
+ return buffer;
+ }
+
public DispatchRequest checkMessageAndReturnSize(java.nio.ByteBuffer
byteBuffer, final boolean checkCRC,
final boolean checkDupInfo) {
return this.checkMessageAndReturnSize(byteBuffer, checkCRC,
checkDupInfo, true);
@@ -456,7 +486,7 @@ public class CommitLog implements Swappable {
}
// 1 TOTAL SIZE
int totalSize = byteBuffer.getInt();
- if (byteBuffer.remaining() < totalSize - 4) {
+ if (totalSize < 0 || byteBuffer.remaining() < totalSize - 4) {
return new DispatchRequest(-1, false /* fail */);
}
@@ -475,7 +505,7 @@ public class CommitLog implements Swappable {
MessageVersion messageVersion =
MessageVersion.valueOfMagicCode(magicCode);
- byte[] bytesContent = new byte[totalSize];
+ byte[] bytesContent = borrowCheckMessageBuffer(totalSize);
int bodyCRC = byteBuffer.getInt();
diff --git
a/store/src/main/java/org/apache/rocketmq/store/config/MessageStoreConfig.java
b/store/src/main/java/org/apache/rocketmq/store/config/MessageStoreConfig.java
index b6a6b6b334..3d1a988e24 100644
---
a/store/src/main/java/org/apache/rocketmq/store/config/MessageStoreConfig.java
+++
b/store/src/main/java/org/apache/rocketmq/store/config/MessageStoreConfig.java
@@ -193,6 +193,15 @@ public class MessageStoreConfig {
// The maximum size of message body,default is 4M,4M only for body
length,not include others.
private int maxMessageSize = 1024 * 1024 * 4;
+ // Upper bound (in bytes) of the per-thread reusable scratch buffer that
+ // CommitLog.checkMessageAndReturnSize keeps for message verification. The
buffer is grow-only,
+ // so this also caps how much memory each dispatch/recovery thread can
retain for its lifetime.
+ // Messages larger than this are verified with a transient buffer instead
of the reusable one.
+ // Default 1M: it covers the vast majority of messages while keeping
per-thread retained memory
+ // small even when concurrent dispatch
(enableBuildConsumeQueueConcurrently) runs many threads;
+ // raise it if larger messages are common.
+ private int maxCheckMessageReuseBufferSize = 1024 * 1024;
+
// The maximum size of message body can be set in config;count with
maxMsgNums * CQ_STORE_UNIT_SIZE(20 || 46)
private int maxFilterMessageSize = 16000;
// Whether check the CRC32 of the records consumed.
@@ -795,6 +804,14 @@ public class MessageStoreConfig {
this.maxMessageSize = maxMessageSize;
}
+ public int getMaxCheckMessageReuseBufferSize() {
+ return maxCheckMessageReuseBufferSize;
+ }
+
+ public void setMaxCheckMessageReuseBufferSize(int
maxCheckMessageReuseBufferSize) {
+ this.maxCheckMessageReuseBufferSize = maxCheckMessageReuseBufferSize;
+ }
+
public int getMaxFilterMessageSize() {
return maxFilterMessageSize;
}
diff --git
a/store/src/test/java/org/apache/rocketmq/store/CheckMessageBufferReuseTest.java
b/store/src/test/java/org/apache/rocketmq/store/CheckMessageBufferReuseTest.java
new file mode 100644
index 0000000000..569979e05f
--- /dev/null
+++
b/store/src/test/java/org/apache/rocketmq/store/CheckMessageBufferReuseTest.java
@@ -0,0 +1,129 @@
+/**
+ * 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;
+
+import java.io.File;
+import java.nio.ByteBuffer;
+import java.util.concurrent.ConcurrentHashMap;
+import org.apache.rocketmq.common.BrokerConfig;
+import org.apache.rocketmq.common.UtilAll;
+import org.apache.rocketmq.store.config.MessageStoreConfig;
+import org.junit.After;
+import org.junit.Before;
+import org.junit.Test;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertNotSame;
+import static org.junit.Assert.assertSame;
+import static org.junit.Assert.assertTrue;
+
+/**
+ * Covers the per-thread reusable scratch buffer used by
+ * {@link CommitLog#checkMessageAndReturnSize} (via {@code
borrowCheckMessageBuffer}):
+ * reuse across messages, grow-only behaviour, the transient fallback for
messages larger than
+ * {@code maxCheckMessageReuseBufferSize}, and rejection of corrupt totalSize
values.
+ */
+public class CheckMessageBufferReuseTest {
+
+ private static final int REUSE_CAP = 64 * 1024;
+ private static final int SMALL = 1024;
+ private static final int LARGER = 8 * 1024;
+ private static final int OVERSIZED = REUSE_CAP + 1;
+
+ private static final String STORE_PATH =
+ System.getProperty("java.io.tmpdir") + File.separator +
"checkbufferreusetest-store";
+
+ private CommitLog commitLog;
+
+ @Before
+ public void init() throws Exception {
+ MessageStoreConfig messageStoreConfig = new MessageStoreConfig();
+ messageStoreConfig.setMappedFileSizeCommitLog(1024 * 8);
+ messageStoreConfig.setMappedFileSizeConsumeQueue(1024 * 4);
+ messageStoreConfig.setMaxHashSlotNum(100);
+ messageStoreConfig.setMaxIndexNum(100 * 10);
+ messageStoreConfig.setMaxCheckMessageReuseBufferSize(REUSE_CAP);
+ messageStoreConfig.setStorePathRootDir(STORE_PATH);
+ messageStoreConfig.setStorePathCommitLog(STORE_PATH + File.separator +
"commitlog");
+ DefaultMessageStore messageStore =
+ new DefaultMessageStore(messageStoreConfig, null, null, new
BrokerConfig(), new ConcurrentHashMap<>());
+ commitLog = new CommitLog(messageStore);
+ }
+
+ @After
+ public void destroy() {
+ UtilAll.deleteFile(new File(STORE_PATH));
+ }
+
+ @Test
+ public void testBufferIsReusedForSameThread() {
+ byte[] first = commitLog.borrowCheckMessageBuffer(SMALL);
+ byte[] second = commitLog.borrowCheckMessageBuffer(SMALL);
+ assertSame("the same thread must reuse the same scratch buffer",
first, second);
+ assertTrue(first.length >= SMALL);
+ }
+
+ @Test
+ public void testBufferGrowsForLargerMessage() {
+ byte[] small = commitLog.borrowCheckMessageBuffer(SMALL);
+ byte[] grown = commitLog.borrowCheckMessageBuffer(LARGER);
+ assertNotSame("a larger message must trigger a grown buffer", small,
grown);
+ assertTrue("grown buffer must fit the larger message", grown.length >=
LARGER);
+ // The grown buffer is retained and reused for subsequent smaller
messages.
+ byte[] again = commitLog.borrowCheckMessageBuffer(SMALL);
+ assertSame("grown buffer must be reused for smaller messages", grown,
again);
+ }
+
+ @Test
+ public void testOversizedMessageUsesTransientBufferAndIsNotRetained() {
+ byte[] reusable = commitLog.borrowCheckMessageBuffer(SMALL);
+
+ byte[] oversized = commitLog.borrowCheckMessageBuffer(OVERSIZED);
+ assertNotSame("oversized message must not use the reusable buffer",
reusable, oversized);
+ assertEquals("oversized transient buffer must fit exactly the
requested size", OVERSIZED, oversized.length);
+
+ // The oversized buffer must not be pinned in the ThreadLocal: a later
small request still
+ // returns the previously cached small buffer, not the oversized one.
+ byte[] afterOversized = commitLog.borrowCheckMessageBuffer(SMALL);
+ assertSame("oversized buffer must not be retained", reusable,
afterOversized);
+ }
+
+ @Test
+ public void testCorruptTotalSizeIsRejectedWithoutAllocating() {
+ // Negative totalSize (corrupt length): must fail before any buffer
allocation, no exception.
+ ByteBuffer negative = ByteBuffer.allocate(16);
+ negative.putInt(-1);
+ negative.putInt(0);
+ negative.putLong(0L);
+ negative.flip();
+ DispatchRequest negativeRequest =
commitLog.checkMessageAndReturnSize(negative, false, false);
+ assertFalse("negative totalSize must be rejected",
negativeRequest.isSuccess());
+ assertEquals(-1, negativeRequest.getMsgSize());
+
+ // totalSize larger than the remaining bytes (truncated/corrupt): must
also fail safely.
+ ByteBuffer tooLarge = ByteBuffer.allocate(16);
+ tooLarge.putInt(1_000_000);
+ tooLarge.putInt(0);
+ tooLarge.putLong(0L);
+ tooLarge.flip();
+ DispatchRequest tooLargeRequest =
commitLog.checkMessageAndReturnSize(tooLarge, false, false);
+ assertFalse("totalSize exceeding remaining bytes must be rejected",
tooLargeRequest.isSuccess());
+ assertEquals(-1, tooLargeRequest.getMsgSize());
+ }
+}