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();
         }
     }
 }

Reply via email to