From 245264ce609c8655828fce64fa2e2319305aa9fd Mon Sep 17 00:00:00 2001 From: Rui <1685901819@qq.com> Date: Wed, 5 Aug 2026 00:51:41 +0800 Subject: [PATCH] [ISSUE #10813] Reduce temporary allocations in the LMQ append path Signed-off-by: Rui <1685901819@qq.com> --- .../org/apache/rocketmq/store/CommitLog.java | 32 ++- .../apache/rocketmq/store/LmqDispatch.java | 55 +++- .../rocketmq/store/LmqDispatchTest.java | 263 ++++++++++++++++++ 3 files changed, 324 insertions(+), 26 deletions(-) create mode 100644 store/src/test/java/org/apache/rocketmq/store/LmqDispatchTest.java 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 260f021aed7..382661800ee 100644 --- a/store/src/main/java/org/apache/rocketmq/store/CommitLog.java +++ b/store/src/main/java/org/apache/rocketmq/store/CommitLog.java @@ -1939,17 +1939,7 @@ public AppendMessageResult handlePropertiesForLmqMsg(ByteBuffer preEncodeBuffer, return null; } - try { - LmqDispatch.wrapLmqDispatch(defaultMessageStore, msgInner); - } catch (ConsumeQueueException e) { - if (e.getCause() instanceof RocksDBException) { - log.error("Failed to wrap multi-dispatch", e); - return new AppendMessageResult(AppendMessageStatus.ROCKSDB_ERROR); - } - log.error("Failed to wrap multi-dispatch", e); - return new AppendMessageResult(AppendMessageStatus.UNKNOWN_ERROR); - } - + LmqDispatch.reinsertWaitStorePropertyForLegacySerialization(msgInner); msgInner.setPropertiesString(MessageDecoder.messageProperties2String(msgInner.getProperties())); final byte[] propertiesData = @@ -2003,7 +1993,20 @@ public AppendMessageResult doAppend(final long fileFromOffset, final ByteBuffer ByteBuffer preEncodeBuffer = msgInner.getEncodedBuff(); boolean isMultiDispatchMsg = messageStoreConfig.isEnableLmq() && msgInner.needDispatchLMQ(); + String[] lmqQueueNames = null; if (isMultiDispatchMsg) { + if (!msgInner.isEncodeCompleted()) { + try { + lmqQueueNames = LmqDispatch.prepareLmqDispatch(defaultMessageStore, msgInner); + } catch (ConsumeQueueException e) { + if (e.getCause() instanceof RocksDBException) { + log.error("Failed to wrap multi-dispatch", e); + return new AppendMessageResult(AppendMessageStatus.ROCKSDB_ERROR); + } + log.error("Failed to wrap multi-dispatch", e); + return new AppendMessageResult(AppendMessageStatus.UNKNOWN_ERROR); + } + } AppendMessageResult appendMessageResult = handlePropertiesForLmqMsg(preEncodeBuffer, msgInner); if (appendMessageResult != null) { return appendMessageResult; @@ -2099,7 +2102,12 @@ public AppendMessageResult doAppend(final long fileFromOffset, final ByteBuffer if (isMultiDispatchMsg) { try { - LmqDispatch.updateLmqOffsets(defaultMessageStore, msgInner); + if (lmqQueueNames == null) { + // The encoded message may be retried after reaching the end of a mapped file. + LmqDispatch.updateLmqOffsets(defaultMessageStore, msgInner); + } else { + LmqDispatch.updateLmqOffsets(defaultMessageStore, lmqQueueNames); + } } catch (ConsumeQueueException e) { // Increase in-memory max offset of the queue should not fail. return new AppendMessageResult(AppendMessageStatus.UNKNOWN_ERROR); diff --git a/store/src/main/java/org/apache/rocketmq/store/LmqDispatch.java b/store/src/main/java/org/apache/rocketmq/store/LmqDispatch.java index 2805f510140..5f589ddab15 100644 --- a/store/src/main/java/org/apache/rocketmq/store/LmqDispatch.java +++ b/store/src/main/java/org/apache/rocketmq/store/LmqDispatch.java @@ -16,7 +16,6 @@ */ package org.apache.rocketmq.store; -import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.common.MixAll; import org.apache.rocketmq.common.message.MessageAccessor; import org.apache.rocketmq.common.message.MessageConst; @@ -28,25 +27,53 @@ public class LmqDispatch { public static void wrapLmqDispatch(MessageStore messageStore, final MessageExtBrokerInner msg) throws ConsumeQueueException { - String lmqNames = msg.getProperty(MessageConst.PROPERTY_INNER_MULTI_DISPATCH); - String[] queueNames = lmqNames.split(MixAll.LMQ_DISPATCH_SEPARATOR); - Long[] queueOffsets = new Long[queueNames.length]; - if (messageStore.getMessageStoreConfig().isEnableLmq()) { - for (int i = 0; i < queueNames.length; i++) { - if (MixAll.isLmq(queueNames[i])) { - queueOffsets[i] = messageStore.getQueueStore().getLmqQueueOffset(queueNames[i], MixAll.LMQ_QUEUE_ID); - } + populateLmqOffsets(messageStore, msg); + msg.removeWaitStorePropertyString(); + } + + static String[] prepareLmqDispatch(MessageStore messageStore, final MessageExtBrokerInner msg) + throws ConsumeQueueException { + return populateLmqOffsets(messageStore, msg); + } + + static void reinsertWaitStorePropertyForLegacySerialization(final MessageExtBrokerInner msg) { + // Reproduce the legacy remove/reinsert mutation without the discarded serialization. + if (msg.getProperties().containsKey(MessageConst.PROPERTY_WAIT_STORE_MSG_OK)) { + String waitStoreMsgOKValue = msg.getProperties().remove(MessageConst.PROPERTY_WAIT_STORE_MSG_OK); + msg.getProperties().put(MessageConst.PROPERTY_WAIT_STORE_MSG_OK, waitStoreMsgOKValue); + } + } + + private static String[] populateLmqOffsets(MessageStore messageStore, final MessageExtBrokerInner msg) + throws ConsumeQueueException { + String[] queueNames = parseLmqQueueNames(msg); + StringBuilder queueOffsets = new StringBuilder(); + boolean enableLmq = messageStore.getMessageStoreConfig().isEnableLmq(); + for (int i = 0; i < queueNames.length; i++) { + if (i > 0) { + queueOffsets.append(MixAll.LMQ_DISPATCH_SEPARATOR); + } + if (enableLmq && MixAll.isLmq(queueNames[i])) { + queueOffsets.append(messageStore.getQueueStore().getLmqQueueOffset(queueNames[i], + MixAll.LMQ_QUEUE_ID)); } } - MessageAccessor.putProperty(msg, MessageConst.PROPERTY_INNER_MULTI_QUEUE_OFFSET, - StringUtils.join(queueOffsets, MixAll.LMQ_DISPATCH_SEPARATOR)); - msg.removeWaitStorePropertyString(); + MessageAccessor.putProperty(msg, MessageConst.PROPERTY_INNER_MULTI_QUEUE_OFFSET, queueOffsets.toString()); + return queueNames; + } + + private static String[] parseLmqQueueNames(final MessageExtBrokerInner msg) { + String lmqNames = msg.getProperty(MessageConst.PROPERTY_INNER_MULTI_DISPATCH); + return lmqNames.split(MixAll.LMQ_DISPATCH_SEPARATOR); } public static void updateLmqOffsets(MessageStore messageStore, final MessageExtBrokerInner msgInner) throws ConsumeQueueException { - String lmqNames = msgInner.getProperty(MessageConst.PROPERTY_INNER_MULTI_DISPATCH); - String[] queueNames = lmqNames.split(MixAll.LMQ_DISPATCH_SEPARATOR); + updateLmqOffsets(messageStore, parseLmqQueueNames(msgInner)); + } + + static void updateLmqOffsets(MessageStore messageStore, String[] queueNames) + throws ConsumeQueueException { for (String queueName : queueNames) { if (messageStore.getMessageStoreConfig().isEnableLmq() && MixAll.isLmq(queueName)) { messageStore.getQueueStore().increaseLmqOffset(queueName, MixAll.LMQ_QUEUE_ID, VALUE_OF_EACH_INCREMENT); diff --git a/store/src/test/java/org/apache/rocketmq/store/LmqDispatchTest.java b/store/src/test/java/org/apache/rocketmq/store/LmqDispatchTest.java new file mode 100644 index 00000000000..421ecbbe0ff --- /dev/null +++ b/store/src/test/java/org/apache/rocketmq/store/LmqDispatchTest.java @@ -0,0 +1,263 @@ +/* + * 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 static org.junit.Assert.assertArrayEquals; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertTrue; +import static org.mockito.ArgumentMatchers.anyInt; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import java.io.File; +import java.net.InetSocketAddress; +import java.nio.ByteBuffer; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import org.apache.rocketmq.common.BrokerConfig; +import org.apache.rocketmq.common.MixAll; +import org.apache.rocketmq.common.UtilAll; +import org.apache.rocketmq.common.message.MessageAccessor; +import org.apache.rocketmq.common.message.MessageConst; +import org.apache.rocketmq.common.message.MessageDecoder; +import org.apache.rocketmq.common.message.MessageExt; +import org.apache.rocketmq.common.message.MessageExtBrokerInner; +import org.apache.rocketmq.store.config.MessageStoreConfig; +import org.apache.rocketmq.store.exception.ConsumeQueueException; +import org.apache.rocketmq.store.queue.ConsumeQueueStoreInterface; +import org.junit.Test; +import org.rocksdb.RocksDBException; + +public class LmqDispatchTest { + + @Test + public void testPrepareAndUpdateMixedQueues() throws Exception { + MessageStore messageStore = mock(MessageStore.class); + MessageStoreConfig messageStoreConfig = mock(MessageStoreConfig.class); + ConsumeQueueStoreInterface queueStore = mock(ConsumeQueueStoreInterface.class); + when(messageStore.getMessageStoreConfig()).thenReturn(messageStoreConfig); + when(messageStore.getQueueStore()).thenReturn(queueStore); + when(messageStoreConfig.isEnableLmq()).thenReturn(true); + + String firstLmq = MixAll.LMQ_PREFIX + "first"; + String secondLmq = MixAll.LMQ_PREFIX + "second"; + when(queueStore.getLmqQueueOffset(firstLmq, MixAll.LMQ_QUEUE_ID)).thenReturn(7L); + when(queueStore.getLmqQueueOffset(secondLmq, MixAll.LMQ_QUEUE_ID)).thenReturn(11L); + + MessageExtBrokerInner message = new MessageExtBrokerInner(); + MessageAccessor.putProperty(message, MessageConst.PROPERTY_INNER_MULTI_DISPATCH, + firstLmq + MixAll.LMQ_DISPATCH_SEPARATOR + "normal-topic" + MixAll.LMQ_DISPATCH_SEPARATOR + secondLmq); + + String[] queueNames = LmqDispatch.prepareLmqDispatch(messageStore, message); + + assertArrayEquals(new String[] {firstLmq, "normal-topic", secondLmq}, queueNames); + assertEquals("7,,11", message.getProperty(MessageConst.PROPERTY_INNER_MULTI_QUEUE_OFFSET)); + verify(queueStore).getLmqQueueOffset(firstLmq, MixAll.LMQ_QUEUE_ID); + verify(queueStore).getLmqQueueOffset(secondLmq, MixAll.LMQ_QUEUE_ID); + + LmqDispatch.updateLmqOffsets(messageStore, queueNames); + verify(queueStore).increaseLmqOffset(firstLmq, MixAll.LMQ_QUEUE_ID, (short) 1); + verify(queueStore).increaseLmqOffset(secondLmq, MixAll.LMQ_QUEUE_ID, (short) 1); + } + + @Test + public void testPublicWrapPreservesWaitPropertyBehaviorWhenLmqIsDisabled() throws Exception { + MessageStore messageStore = mock(MessageStore.class); + MessageStoreConfig messageStoreConfig = mock(MessageStoreConfig.class); + ConsumeQueueStoreInterface queueStore = mock(ConsumeQueueStoreInterface.class); + when(messageStore.getMessageStoreConfig()).thenReturn(messageStoreConfig); + when(messageStore.getQueueStore()).thenReturn(queueStore); + when(messageStoreConfig.isEnableLmq()).thenReturn(false); + + MessageExtBrokerInner message = new MessageExtBrokerInner(); + message.setWaitStoreMsgOK(true); + MessageAccessor.putProperty(message, MessageConst.PROPERTY_INNER_MULTI_DISPATCH, + MixAll.LMQ_PREFIX + "first,normal-topic," + MixAll.LMQ_PREFIX + "second"); + + LmqDispatch.wrapLmqDispatch(messageStore, message); + + assertEquals(",,", message.getProperty(MessageConst.PROPERTY_INNER_MULTI_QUEUE_OFFSET)); + assertEquals("true", message.getProperty(MessageConst.PROPERTY_WAIT_STORE_MSG_OK)); + Map serializedProperties = + MessageDecoder.string2messageProperties(message.getPropertiesString()); + assertFalse(serializedProperties.containsKey(MessageConst.PROPERTY_WAIT_STORE_MSG_OK)); + verify(queueStore, never()).getLmqQueueOffset(anyString(), anyInt()); + } + + @Test + public void testCommitLogPreservesLegacyPropertyBytesAndIncrementsOffsetOnce() throws Exception { + AppendFixture fixture = createAppendFixture(); + try { + String lmqName = MixAll.LMQ_PREFIX + "put-ok"; + MessageExtBrokerInner message = createLmqMessage(fixture.messageStoreConfig, lmqName); + MessageExtBrokerInner legacyMessage = createLmqMessage(fixture.messageStoreConfig, lmqName); + MessageAccessor.putProperty(legacyMessage, MessageConst.PROPERTY_INNER_MULTI_QUEUE_OFFSET, "0"); + legacyMessage.removeWaitStorePropertyString(); + legacyMessage.setPropertiesString(MessageDecoder.messageProperties2String(legacyMessage.getProperties())); + byte[] expectedProperties = legacyMessage.getPropertiesString().getBytes(MessageDecoder.CHARSET_UTF8); + ByteBuffer destination = ByteBuffer.allocate(1024); + + AppendMessageResult result = fixture.callback.doAppend(0, destination, destination.capacity(), message, + null); + + assertEquals(AppendMessageStatus.PUT_OK, result.getStatus()); + assertEquals(1L, fixture.messageStore.getQueueStore().getLmqQueueOffset(lmqName, MixAll.LMQ_QUEUE_ID)); + assertEquals(legacyMessage.getPropertiesString(), message.getPropertiesString()); + + int messageLength = destination.getInt(0); + byte[] persistedProperties = new byte[expectedProperties.length]; + ByteBuffer persistedPropertiesBuffer = destination.duplicate(); + persistedPropertiesBuffer.position(messageLength - expectedProperties.length); + persistedPropertiesBuffer.get(persistedProperties); + assertArrayEquals(expectedProperties, persistedProperties); + + MessageExt persistedMessage = MessageDecoder.decode((ByteBuffer) destination.flip()); + assertNotNull(persistedMessage); + assertEquals("true", persistedMessage.getProperty(MessageConst.PROPERTY_WAIT_STORE_MSG_OK)); + assertEquals("0", persistedMessage.getProperty(MessageConst.PROPERTY_INNER_MULTI_QUEUE_OFFSET)); + } finally { + fixture.destroy(); + } + } + + @Test + public void testCommitLogEndOfFileRetryIncrementsOffsetOnce() throws Exception { + AppendFixture fixture = createAppendFixture(); + try { + String lmqName = MixAll.LMQ_PREFIX + "end-of-file"; + MessageExtBrokerInner message = createLmqMessage(fixture.messageStoreConfig, lmqName); + PutMessageContext putMessageContext = new PutMessageContext("test-topic-0"); + ByteBuffer endOfFileBuffer = ByteBuffer.allocate(8); + + AppendMessageResult endOfFileResult = fixture.callback.doAppend(0, endOfFileBuffer, + endOfFileBuffer.capacity(), message, putMessageContext); + + assertEquals(AppendMessageStatus.END_OF_FILE, endOfFileResult.getStatus()); + assertTrue(message.isEncodeCompleted()); + assertEquals(0L, fixture.messageStore.getQueueStore().getLmqQueueOffset(lmqName, + MixAll.LMQ_QUEUE_ID)); + + ByteBuffer destination = ByteBuffer.allocate(1024); + AppendMessageResult putResult = fixture.callback.doAppend(8, destination, destination.capacity(), message, + putMessageContext); + + assertEquals(AppendMessageStatus.PUT_OK, putResult.getStatus()); + assertEquals(1L, fixture.messageStore.getQueueStore().getLmqQueueOffset(lmqName, MixAll.LMQ_QUEUE_ID)); + } finally { + fixture.destroy(); + } + } + + @Test + public void testCommitLogMapsRocksDbAndConsumeQueueFailures() throws Exception { + assertPrepareFailureStatus(new ConsumeQueueException(new RocksDBException("rocksdb failure")), + AppendMessageStatus.ROCKSDB_ERROR); + assertPrepareFailureStatus(new ConsumeQueueException("consume queue failure"), + AppendMessageStatus.UNKNOWN_ERROR); + } + + private void assertPrepareFailureStatus(ConsumeQueueException exception, AppendMessageStatus expectedStatus) + throws Exception { + String storePath = newStorePath(); + MessageStoreConfig messageStoreConfig = createMessageStoreConfig(storePath); + DefaultMessageStore messageStore = mock(DefaultMessageStore.class); + ConsumeQueueStoreInterface queueStore = mock(ConsumeQueueStoreInterface.class); + when(messageStore.getMessageStoreConfig()).thenReturn(messageStoreConfig); + when(messageStore.getQueueStore()).thenReturn(queueStore); + when(queueStore.getLmqQueueOffset(anyString(), anyInt())).thenThrow(exception); + + CommitLog commitLog = new CommitLog(messageStore); + AppendMessageCallback callback = commitLog.new DefaultAppendMessageCallback(messageStoreConfig); + MessageExtBrokerInner message = createLmqMessage(messageStoreConfig, MixAll.LMQ_PREFIX + "failure"); + + AppendMessageResult result = callback.doAppend(0, ByteBuffer.allocate(1024), 1024, message, null); + + assertEquals(expectedStatus, result.getStatus()); + assertFalse(message.isEncodeCompleted()); + UtilAll.deleteFile(new File(storePath)); + } + + private MessageExtBrokerInner createLmqMessage(MessageStoreConfig messageStoreConfig, String lmqName) { + MessageExtBrokerInner message = new MessageExtBrokerInner(); + message.setTopic("test-topic"); + message.setQueueId(0); + message.setBody("body".getBytes(MessageDecoder.CHARSET_UTF8)); + message.setBornTimestamp(System.currentTimeMillis()); + message.setStoreTimestamp(System.currentTimeMillis()); + message.setBornHost(new InetSocketAddress("127.0.0.1", 12345)); + message.setStoreHost(new InetSocketAddress("127.0.0.1", 10911)); + message.setWaitStoreMsgOK(true); + message.putUserProperty("m", "same-bucket-as-wait"); + MessageAccessor.putProperty(message, MessageConst.PROPERTY_INNER_MULTI_DISPATCH, lmqName); + message.setPropertiesString(MessageDecoder.messageProperties2String(message.getProperties())); + + MessageExtEncoder encoder = new MessageExtEncoder(messageStoreConfig); + assertNull(encoder.encode(message)); + message.setEncodedBuff(encoder.getEncoderBuffer()); + return message; + } + + private AppendFixture createAppendFixture() throws Exception { + String storePath = newStorePath(); + MessageStoreConfig messageStoreConfig = createMessageStoreConfig(storePath); + DefaultMessageStore messageStore = new DefaultMessageStore(messageStoreConfig, null, null, + new BrokerConfig(), new ConcurrentHashMap<>()); + CommitLog commitLog = new CommitLog(messageStore); + return new AppendFixture(storePath, messageStoreConfig, messageStore, + commitLog.new DefaultAppendMessageCallback(messageStoreConfig)); + } + + private MessageStoreConfig createMessageStoreConfig(String storePath) { + MessageStoreConfig messageStoreConfig = new MessageStoreConfig(); + messageStoreConfig.setStorePathRootDir(storePath); + messageStoreConfig.setStorePathCommitLog(storePath + File.separator + "commitlog"); + messageStoreConfig.setMappedFileSizeCommitLog(8 * 1024); + messageStoreConfig.setMaxMessageSize(1024 * 1024); + messageStoreConfig.setEnableLmq(true); + return messageStoreConfig; + } + + private String newStorePath() { + return System.getProperty("java.io.tmpdir") + File.separator + "lmq-dispatch-" + UUID.randomUUID(); + } + + private static class AppendFixture { + private final String storePath; + private final MessageStoreConfig messageStoreConfig; + private final DefaultMessageStore messageStore; + private final AppendMessageCallback callback; + + private AppendFixture(String storePath, MessageStoreConfig messageStoreConfig, + DefaultMessageStore messageStore, AppendMessageCallback callback) { + this.storePath = storePath; + this.messageStoreConfig = messageStoreConfig; + this.messageStore = messageStore; + this.callback = callback; + } + + private void destroy() { + UtilAll.deleteFile(new File(storePath)); + } + } +}