This is an automated email from the ASF dual-hosted git repository.
lollipopjin 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 050d2b42b0 [ISSUE #11168] Fix tiered storage commit failure reporting
and related hazards (#11169)
050d2b42b0 is described below
commit 050d2b42b08bdfdf3ed273084d4aee2a9be7135c
Author: lizhimins <[email protected]>
AuthorDate: Sun Sep 20 11:35:07 2026 +0800
[ISSUE #11168] Fix tiered storage commit failure reporting and related
hazards (#11169)
---
.../tieredstore/core/MessageStoreFetcherImpl.java | 7 +++-
.../exception/TieredStoreErrorCode.java | 7 ++++
.../exception/TieredStoreException.java | 10 +++++
.../rocketmq/tieredstore/index/IndexStoreFile.java | 21 ++++++++++-
.../rocketmq/tieredstore/provider/FileSegment.java | 44 ++++++++++++----------
.../core/MessageStoreFetcherImplTest.java | 37 ++++++++++++++++++
.../exception/TieredStoreExceptionTest.java | 14 +++++++
.../tieredstore/provider/FileSegmentTest.java | 5 ++-
8 files changed, 122 insertions(+), 23 deletions(-)
diff --git
a/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/core/MessageStoreFetcherImpl.java
b/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/core/MessageStoreFetcherImpl.java
index 84f0c359ce..f4590c0c7a 100644
---
a/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/core/MessageStoreFetcherImpl.java
+++
b/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/core/MessageStoreFetcherImpl.java
@@ -459,7 +459,12 @@ public class MessageStoreFetcherImpl implements
MessageStoreFetcher {
return CompletableFuture.completedFuture(result);
}
- boolean cacheBusy = fetcherCache.estimatedSize() > memoryMaxSize * 0.8;
+ // The cache is bounded by maximumWeight (bytes, via the
SelectBufferResult#getSize weigher),
+ // so compare against weightedSize() rather than estimatedSize(),
which counts entries.
+ long cacheWeight = fetcherCache.policy().eviction()
+ .map(eviction -> eviction.weightedSize().orElse(0L))
+ .orElse(0L);
+ boolean cacheBusy = cacheWeight > memoryMaxSize * 0.8;
if (storeConfig.isReadAheadCacheEnable() && !cacheBusy) {
return getMessageFromCacheAsync(flatFile, group, queueOffset,
maxCount, messageFilter);
} else {
diff --git
a/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/exception/TieredStoreErrorCode.java
b/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/exception/TieredStoreErrorCode.java
index d29025f1c5..8afd506b1d 100644
---
a/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/exception/TieredStoreErrorCode.java
+++
b/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/exception/TieredStoreErrorCode.java
@@ -53,6 +53,13 @@ public enum TieredStoreErrorCode {
*/
SEGMENT_SEALED,
+ /**
+ * Error code for an object that does not exist in the storage system. A
caller that can still
+ * answer from the remaining segments should treat this as an empty result
rather than a query
+ * failure, since the object may be deleted while a read is already in
flight.
+ */
+ FILE_NOT_FOUND,
+
/**
* Error code for an unknown error.
*/
diff --git
a/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/exception/TieredStoreException.java
b/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/exception/TieredStoreException.java
index 3841643299..483108946b 100644
---
a/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/exception/TieredStoreException.java
+++
b/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/exception/TieredStoreException.java
@@ -27,6 +27,16 @@ public class TieredStoreException extends RuntimeException {
this.errorCode = errorCode;
}
+ public static boolean hasErrorCode(Throwable throwable,
TieredStoreErrorCode errorCode) {
+ for (Throwable cause = throwable; cause != null && cause !=
cause.getCause(); cause = cause.getCause()) {
+ if (cause instanceof TieredStoreException &&
+ errorCode == ((TieredStoreException) cause).getErrorCode()) {
+ return true;
+ }
+ }
+ return false;
+ }
+
public TieredStoreErrorCode getErrorCode() {
return errorCode;
}
diff --git
a/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/index/IndexStoreFile.java
b/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/index/IndexStoreFile.java
index 8fd4b2961b..60d92b0c68 100644
---
a/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/index/IndexStoreFile.java
+++
b/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/index/IndexStoreFile.java
@@ -29,6 +29,7 @@ import java.util.List;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CompletionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
@@ -41,6 +42,8 @@ import org.apache.rocketmq.store.logfile.DefaultMappedFile;
import org.apache.rocketmq.store.logfile.MappedFile;
import org.apache.rocketmq.tieredstore.MessageStoreConfig;
import org.apache.rocketmq.tieredstore.common.AppendResult;
+import org.apache.rocketmq.tieredstore.exception.TieredStoreErrorCode;
+import org.apache.rocketmq.tieredstore.exception.TieredStoreException;
import org.apache.rocketmq.tieredstore.provider.FileSegment;
import org.apache.rocketmq.tieredstore.provider.PosixFileSegment;
import org.apache.rocketmq.tieredstore.util.MessageStoreUtil;
@@ -418,8 +421,14 @@ public class IndexStoreFile implements IndexFile {
return future.whenComplete((result, throwable) -> {
long costTime = stopwatch.elapsed(TimeUnit.MILLISECONDS);
if (throwable != null) {
- log.error("IndexStoreFile#queryAsyncFromSegmentFile, query
from segment file error, cost={}ms, timestamp={}, key={}, hashCode={},
maxCount={}, timeRange={}-{}",
- costTime, getTimestamp(), key, hashCode, maxCount,
beginTime, endTime, throwable);
+ if (TieredStoreException.hasErrorCode(throwable,
TieredStoreErrorCode.FILE_NOT_FOUND)) {
+ log.info("IndexStoreFile#queryAsyncFromSegmentFile,
segment file not found, treat as no result, cost={}ms, timestamp={}, key={},
hashCode={}, maxCount={}, timeRange={}-{}, reason={}",
+ costTime, getTimestamp(), key, hashCode, maxCount,
beginTime, endTime, throwable.getMessage());
+ } else {
+ // The exception propagates to IndexStoreService, which
records it at ERROR.
+ log.debug("IndexStoreFile#queryAsyncFromSegmentFile, query
from segment file error, cost={}ms, timestamp={}, key={}, hashCode={},
maxCount={}, timeRange={}-{}",
+ costTime, getTimestamp(), key, hashCode, maxCount,
beginTime, endTime, throwable);
+ }
} else {
String details = Optional.ofNullable(result)
.map(r -> r.stream()
@@ -430,6 +439,14 @@ public class IndexStoreFile implements IndexFile {
log.debug("IndexStoreFile#queryAsyncFromSegmentFile, query
from segment file, cost={}ms, timestamp={}, resultSize={}, ({}), key={},
hashCode={}, maxCount={}, timeRange={}-{}",
costTime, getTimestamp(), result != null ? result.size() :
0, details, key, hashCode, maxCount, beginTime, endTime);
}
+ }).exceptionally(throwable -> {
+ if (!TieredStoreException.hasErrorCode(throwable,
TieredStoreErrorCode.FILE_NOT_FOUND)) {
+ if (throwable instanceof RuntimeException) {
+ throw (RuntimeException) throwable;
+ }
+ throw new CompletionException(throwable);
+ }
+ return Collections.emptyList();
});
}
diff --git
a/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/provider/FileSegment.java
b/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/provider/FileSegment.java
index 0d4e39b74f..cc659c7c03 100644
---
a/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/provider/FileSegment.java
+++
b/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/provider/FileSegment.java
@@ -237,8 +237,11 @@ public abstract class FileSegment implements
Comparable<FileSegment>, FileSegmen
if (fileSegmentInputStream != null) {
long fileSize = this.getSize();
if (fileSize == GET_FILE_SIZE_ERROR) {
- log.error("FileSegment#commitAsync, correct position error,
fileName={}, commit={}, append={}, buffer={}",
- this.getPath(), commitPosition, appendPosition,
fileSegmentInputStream.getContentLength());
+ long contentLength = fileSegmentInputStream.getContentLength();
+ log.error("FileSegment#commitAsync, fileName={}, result={},
commit={}, content={}, " +
+ "expect={}, append={}, remote={}",
+ this.getPath(), "SIZE_LOOKUP_FAILED", commitPosition,
contentLength,
+ commitPosition + contentLength, appendPosition, fileSize);
releaseCommitLock();
return CompletableFuture.completedFuture(false);
}
@@ -280,29 +283,32 @@ public abstract class FileSegment implements
Comparable<FileSegment>, FileSegmen
}
private boolean handleCommitException(Throwable e) {
-
- log.warn("FileSegment#handleCommitException, commit exception,
filePath={}", this.filePath, e);
-
- // Get root cause here
Throwable rootCause = e.getCause() != null ? e.getCause() : e;
+ long commitPositionBefore = commitPosition;
+ long contentLength = fileSegmentInputStream.getContentLength();
+ long expectPosition = commitPositionBefore + contentLength;
long fileSize = rootCause instanceof TieredStoreException ?
- ((TieredStoreException) rootCause).getPosition() : this.getSize();
-
- long expectPosition = commitPosition +
fileSegmentInputStream.getContentLength();
- if (fileSize == GET_FILE_SIZE_ERROR) {
- log.error("FileSegment#handleCommitException, get file size error
after commit, fileName={}, commit={}, content={}, expect={}, append={}",
- this.getPath(), commitPosition,
fileSegmentInputStream.getContentLength(), expectPosition, appendPosition);
- return false;
- }
-
- if (correctPosition(fileSize)) {
+ ((TieredStoreException) rootCause).getPosition() :
GET_FILE_SIZE_ERROR;
+ boolean sizeKnown = fileSize != GET_FILE_SIZE_ERROR;
+
+ boolean landed = false;
+ String result;
+ if (!sizeKnown) {
+ result = "RETRY_AFTER_RECONCILE";
+ } else if (correctPosition(fileSize)) {
fileSegmentInputStream = null;
- return true;
+ result = "REMOTE_LANDED";
+ landed = true;
} else {
fileSegmentInputStream.rewind();
- return false;
+ result = "RETRY_AFTER_REWIND";
}
+
+ log.warn("FileSegment#handleCommitException, fileName={}, result={},
commit={}, content={}, " +
+ "expect={}, append={}, remote={}",
+ this.getPath(), result, commitPositionBefore, contentLength,
expectPosition, appendPosition, fileSize, e);
+ return landed;
}
private void releaseCommitLock() {
@@ -352,7 +358,7 @@ public abstract class FileSegment implements
Comparable<FileSegment>, FileSegmen
int readableBytes = (int) (currentCommitPosition - position);
if (readableBytes < length) {
- log.debug("FileSegment#readAsync, request position exceeds commit
position, " +
+ log.warn("FileSegment#readAsync, request position exceeds commit
position, " +
"file={}, requestPosition={}, commitPosition={},
changeLength={} to {}",
getPath(), position, currentCommitPosition, length,
readableBytes);
length = readableBytes;
diff --git
a/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/core/MessageStoreFetcherImplTest.java
b/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/core/MessageStoreFetcherImplTest.java
index fd681f27b7..b5c39c5252 100644
---
a/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/core/MessageStoreFetcherImplTest.java
+++
b/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/core/MessageStoreFetcherImplTest.java
@@ -20,6 +20,7 @@ import com.google.common.collect.Sets;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.time.Duration;
+import java.util.concurrent.CompletableFuture;
import java.util.concurrent.atomic.AtomicLong;
import org.apache.commons.lang3.reflect.FieldUtils;
import org.apache.rocketmq.common.BoundaryType;
@@ -33,8 +34,11 @@ import org.apache.rocketmq.store.MessageFilter;
import org.apache.rocketmq.store.QueryMessageResult;
import org.apache.rocketmq.tieredstore.MessageStoreConfig;
import org.apache.rocketmq.tieredstore.TieredMessageStore;
+import org.apache.rocketmq.tieredstore.common.GetMessageResultExt;
import org.apache.rocketmq.tieredstore.common.SelectBufferResult;
+import org.apache.rocketmq.tieredstore.file.FlatFileStore;
import org.apache.rocketmq.tieredstore.file.FlatMessageFile;
+import org.apache.rocketmq.tieredstore.index.IndexService;
import org.apache.rocketmq.tieredstore.util.MessageFormatUtilTest;
import org.apache.rocketmq.tieredstore.util.MessageStoreUtilTest;
import org.awaitility.Awaitility;
@@ -187,6 +191,39 @@ public class MessageStoreFetcherImplTest {
Assert.assertEquals(100 / times.get(), batchSize);
}
+ @Test
+ public void cacheWeightControlsReadPathTest() {
+ MessageStoreConfig config = new MessageStoreConfig();
+ config.setReadAheadCacheEnable(true);
+ config.setReadAheadCacheSizeThresholdRate(1024D /
Runtime.getRuntime().maxMemory());
+
+ TieredMessageStore tieredStore =
Mockito.mock(TieredMessageStore.class);
+ FlatFileStore flatFileStore = Mockito.mock(FlatFileStore.class);
+ FlatMessageFile flatFile = Mockito.mock(FlatMessageFile.class);
+
Mockito.when(flatFileStore.getFlatFile(Mockito.any(MessageQueue.class))).thenReturn(flatFile);
+ Mockito.when(flatFile.getConsumeQueueMinOffset()).thenReturn(0L);
+ Mockito.when(flatFile.getConsumeQueueCommitOffset()).thenReturn(100L);
+
+ MessageStoreFetcherImpl cacheFetcher = Mockito.spy(new
MessageStoreFetcherImpl(
+ tieredStore, config, flatFileStore,
Mockito.mock(IndexService.class)));
+ Mockito.doReturn(CompletableFuture.completedFuture(new
GetMessageResult()))
+ .when(cacheFetcher).getMessageFromCacheAsync(flatFile, groupName,
1L, 1, null);
+ Mockito.doReturn(CompletableFuture.completedFuture(new
GetMessageResultExt()))
+ .when(cacheFetcher).getMessageFromTieredStoreAsync(flatFile, 1L,
1);
+
+ cacheFetcher.getMessageAsync(groupName, "topic", 0, 1L, 1,
null).join();
+ Mockito.verify(cacheFetcher).getMessageFromCacheAsync(flatFile,
groupName, 1L, 1, null);
+ Mockito.verify(cacheFetcher,
Mockito.never()).getMessageFromTieredStoreAsync(flatFile, 1L, 1);
+
+ int entrySize = (int) Math.ceil(cacheFetcher.memoryMaxSize * 0.9);
+ cacheFetcher.getFetcherCache().put("entry", new SelectBufferResult(
+ ByteBuffer.allocate(entrySize), 0, entrySize, 0));
+ cacheFetcher.getFetcherCache().cleanUp();
+
+ cacheFetcher.getMessageAsync(groupName, "topic", 0, 1L, 1,
null).join();
+ Mockito.verify(cacheFetcher).getMessageFromTieredStoreAsync(flatFile,
1L, 1);
+ }
+
@Test
public void getMessageFromCacheTagFilterTest() throws Exception {
dispatcherTest.dispatchFromCommitLogTest();
diff --git
a/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/exception/TieredStoreExceptionTest.java
b/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/exception/TieredStoreExceptionTest.java
index 1de891a8ac..1e6481ae3f 100644
---
a/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/exception/TieredStoreExceptionTest.java
+++
b/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/exception/TieredStoreExceptionTest.java
@@ -16,6 +16,7 @@
*/
package org.apache.rocketmq.tieredstore.exception;
+import java.util.concurrent.CompletionException;
import org.junit.Assert;
import org.junit.Test;
@@ -38,4 +39,17 @@ public class TieredStoreExceptionTest {
Assert.assertEquals(position, tieredStoreException.getPosition());
Assert.assertNotNull(tieredStoreException.toString());
}
+
+ @Test
+ public void hasErrorCodeTest() {
+ Throwable throwable = new CompletionException(
+ new TieredStoreException(TieredStoreErrorCode.FILE_NOT_FOUND, "not
found"));
+
+ Assert.assertTrue(TieredStoreException.hasErrorCode(
+ throwable, TieredStoreErrorCode.FILE_NOT_FOUND));
+ Assert.assertFalse(TieredStoreException.hasErrorCode(
+ throwable, TieredStoreErrorCode.IO_ERROR));
+ Assert.assertFalse(TieredStoreException.hasErrorCode(
+ null, TieredStoreErrorCode.FILE_NOT_FOUND));
+ }
}
\ No newline at end of file
diff --git
a/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/provider/FileSegmentTest.java
b/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/provider/FileSegmentTest.java
index 26844113cd..18091441f8 100644
---
a/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/provider/FileSegmentTest.java
+++
b/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/provider/FileSegmentTest.java
@@ -462,8 +462,11 @@ public class FileSegmentTest {
.thenReturn(CompletableFuture.supplyAsync(() -> {
throw new RuntimeException("Runtime Error for Test");
}));
- Mockito.when(fileSpySegment.getSize()).thenReturn(0L);
Assert.assertFalse(fileSpySegment.commitAsync().join());
+ // handleCommitException runs on the thread that completed the
future, which is a netty IO
+ // thread for a network provider, so it must not do a remote size
lookup there. An unknown
+ // length is reconciled by the next commitAsync instead.
+ Mockito.verify(fileSpySegment, Mockito.never()).getSize();
}
}
}