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

yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/branch-4.1 by this push:
     new cb239b7712b branch-4.1: [fix](cloud) Invalidate version caches on 
visible commit retries (#67813) (#68650)
cb239b7712b is described below

commit cb239b7712b0695a2c2f6fa3ddc9c93d9978012c
Author: Luwei <[email protected]>
AuthorDate: Wed Sep 30 20:55:45 2026 +0800

    branch-4.1: [fix](cloud) Invalidate version caches on visible commit 
retries (#67813) (#68650)
    
    ### What problem does this PR solve?
    
    Issue Number: #67099
    
    Related PR: #67813
    
    Problem Summary: A retry after a lost cloud commit response can return a
    VISIBLE transaction without version results. Invalidate the affected
    table and partition version caches, and notify other FEs so subsequent
    reads fetch current versions from the MetaService. This cherry-picks the
    merged #67813 commit (`3dbe4d3ca53307f73db96eda23d4efc9f02963c5`) to
    `branch-4.1`.
    
    The conflict resolution keeps the branch-4.1 partition cache API, which
    has no commit-TSO argument, and adapts the tests to its Java 8 and JUnit
    4 setup. It preserves the source PR's cache invalidation epochs and
    batch publication locks.
    
    ### Release note
    
    Prevent stale reads and false commit failures after cloud commit
    retries.
    
    ### Check List (For Author)
    
    - Test
    - [x] Unit Test: 51 passed (`CloudGlobalTransactionMgrTest`,
    `CloudPartitionTest`, `OlapTableTest`, `VersionHelperTest`).
        - [x] FE Checkstyle: 0 violations; `git diff --check` passed.
        - [ ] Live cluster / SQL regression: not run.
    - Behavior changed:
    - [x] Yes. VISIBLE retries without versions invalidate local and peer
    caches; subsequent reads fetch fresh versions without changing the
    successful commit outcome.
    - Does this need documentation?
        - [x] No.
---
 .../java/org/apache/doris/catalog/OlapTable.java   |  21 +-
 .../cloud/catalog/CloudFEVersionSynchronizer.java  |  59 +++
 .../apache/doris/cloud/catalog/CloudPartition.java |  43 ++-
 .../transaction/CloudGlobalTransactionMgr.java     |  11 +
 .../transaction/CloudGlobalTransactionMgrTest.java | 400 ++++++++++++++++++++-
 5 files changed, 521 insertions(+), 13 deletions(-)

diff --git a/fe/fe-core/src/main/java/org/apache/doris/catalog/OlapTable.java 
b/fe/fe-core/src/main/java/org/apache/doris/catalog/OlapTable.java
index f9ac99263d2..f319735e1d1 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/catalog/OlapTable.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/catalog/OlapTable.java
@@ -127,6 +127,7 @@ import java.util.Set;
 import java.util.TreeMap;
 import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.atomic.AtomicLong;
 import java.util.concurrent.locks.ReadWriteLock;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 import java.util.stream.Collectors;
@@ -243,6 +244,10 @@ public class OlapTable extends Table implements 
MTMVRelatedTableIf, GsonPostProc
     private volatile long lastTableVersionCachedTimeMs = 0;
     private volatile long cachedTableVersion = -1;
 
+    // Commit notifications cannot clear invalidation: their versions may 
precede the missing commit result.
+    private final AtomicLong tableVersionCacheEpoch = new AtomicLong();
+    private final AtomicLong refreshedTableVersionCacheEpoch = new 
AtomicLong();
+
     private ReadWriteLock versionLock = Config.isCloudMode() ? new 
ReentrantReadWriteLock(true) : null;
 
     public OlapTable() {
@@ -3398,7 +3403,7 @@ public class OlapTable extends Table implements 
MTMVRelatedTableIf, GsonPostProc
     @VisibleForTesting
     protected boolean isCachedTableVersionExpired() {
         // -1 means no cache yet, need to fetch from MS
-        if (cachedTableVersion == -1) {
+        if (cachedTableVersion == -1 || tableVersionCacheEpoch.get() != 
refreshedTableVersionCacheEpoch.get()) {
             return true;
         }
         ConnectContext ctx = ConnectContext.get();
@@ -3412,13 +3417,18 @@ public class OlapTable extends Table implements 
MTMVRelatedTableIf, GsonPostProc
 
     public boolean isCachedTableVersionExpired(long expirationMs) {
         // -1 means no cache yet, need to fetch from MS
-        if (cachedTableVersion == -1 || expirationMs <= 0) {
+        if (cachedTableVersion == -1 || expirationMs <= 0
+                || tableVersionCacheEpoch.get() != 
refreshedTableVersionCacheEpoch.get()) {
             return true;
         }
         return System.currentTimeMillis() - lastTableVersionCachedTimeMs > 
expirationMs;
     }
 
-    public void setCachedTableVersion(long version) {
+    public void invalidateCachedTableVersion() {
+        tableVersionCacheEpoch.incrementAndGet();
+    }
+
+    public synchronized void setCachedTableVersion(long version) {
         if (version >= cachedTableVersion) {
             cachedTableVersion = version;
             lastTableVersionCachedTimeMs = System.currentTimeMillis();
@@ -3439,6 +3449,7 @@ public class OlapTable extends Table implements 
MTMVRelatedTableIf, GsonPostProc
             return getCachedTableVersion();
         }
 
+        long cacheEpoch = tableVersionCacheEpoch.get();
         // get version rpc
         Cloud.GetVersionRequest request = Cloud.GetVersionRequest.newBuilder()
                 .setRequestIp(FrontendOptions.getLocalHostAddressCached())
@@ -3465,6 +3476,7 @@ public class OlapTable extends Table implements 
MTMVRelatedTableIf, GsonPostProc
             }
             // update cache
             setCachedTableVersion(version);
+            refreshedTableVersionCacheEpoch.accumulateAndGet(cacheEpoch, 
Math::max);
             return version;
         } catch (RpcException e) {
             LOG.warn("get version from meta service failed", e);
@@ -3529,9 +3541,11 @@ public class OlapTable extends Table implements 
MTMVRelatedTableIf, GsonPostProc
     private static List<Long> getVisibleVersionInBatchFromMs(List<OlapTable> 
tables) {
         List<Long> dbIds = new ArrayList<>(tables.size());
         List<Long> tableIds = new ArrayList<>(tables.size());
+        List<Long> cacheEpochs = new ArrayList<>(tables.size());
         for (OlapTable table : tables) {
             dbIds.add(table.getDatabase().getId());
             tableIds.add(table.getId());
+            cacheEpochs.add(table.tableVersionCacheEpoch.get());
         }
 
         List<Long> versions = getVisibleVersionFromMeta(dbIds, tableIds);
@@ -3540,6 +3554,7 @@ public class OlapTable extends Table implements 
MTMVRelatedTableIf, GsonPostProc
         Preconditions.checkState(tables.size() == versions.size());
         for (int i = 0; i < tables.size(); i++) {
             tables.get(i).setCachedTableVersion(versions.get(i));
+            
tables.get(i).refreshedTableVersionCacheEpoch.accumulateAndGet(cacheEpochs.get(i),
 Math::max);
         }
         return versions;
     }
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudFEVersionSynchronizer.java
 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudFEVersionSynchronizer.java
index f843b8b58e7..d88348b704a 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudFEVersionSynchronizer.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudFEVersionSynchronizer.java
@@ -49,6 +49,9 @@ import java.util.stream.Collectors;
 
 public class CloudFEVersionSynchronizer {
     private static final Logger LOG = 
LogManager.getLogger(CloudFEVersionSynchronizer.class);
+    // Reuse the version RPC for whole-table invalidation when a commit 
response has no versions.
+    // Older FEs ignore this value and continue relying on their periodic 
version synchronization.
+    private static final long INVALIDATE_VERSION_CACHE = -1;
 
     private static final ExecutorService SYNC_VERSION_THREAD_POOL = 
Executors.newFixedThreadPool(
             Config.cloud_sync_version_task_threads_num,
@@ -57,6 +60,54 @@ public class CloudFEVersionSynchronizer {
     public CloudFEVersionSynchronizer() {
     }
 
+    public void invalidateVersionCaches(long dbId, List<Long> tableIds) {
+        Database db = Env.getCurrentInternalCatalog().getDbNullable(dbId);
+        if (db == null) {
+            return; // The database may have been dropped since the original 
commit.
+        }
+        List<OlapTable> tables = new ArrayList<>();
+        for (long tableId : tableIds) {
+            Table table = db.getTableNullable(tableId);
+            if (table != null && table.isManagedTable()) {
+                tables.add((OlapTable) table);
+            }
+        }
+        invalidateVersionCaches(tables);
+        List<Pair<OlapTable, Long>> tableVersions = tables.stream()
+                .map(table -> Pair.of(table, 
INVALIDATE_VERSION_CACHE)).collect(Collectors.toList());
+        pushVersionAsync(dbId, tableVersions, Collections.emptyMap());
+    }
+
+    private void invalidateVersionCaches(List<OlapTable> tables) {
+        tables.sort(Comparator.comparingLong(OlapTable::getId));
+        for (OlapTable table : tables) {
+            table.readLock();
+        }
+        try {
+            for (OlapTable table : tables) {
+                table.versionWriteLock();
+            }
+            try {
+                for (OlapTable table : tables) {
+                    table.invalidateCachedTableVersion();
+                    for (Partition partition : table.getAllPartitions()) {
+                        ((CloudPartition) 
partition).invalidateCachedVisibleVersion();
+                    }
+                }
+            } finally {
+                for (int i = tables.size() - 1; i >= 0; i--) {
+                    tables.get(i).versionWriteUnlock();
+                }
+            }
+        } finally {
+            for (int i = tables.size() - 1; i >= 0; i--) {
+                tables.get(i).readUnlock();
+            }
+        }
+        LOG.info("Invalidated cloud version caches for tables {}",
+                
tables.stream().map(OlapTable::getId).collect(Collectors.toList()));
+    }
+
     // master FE send sync version rpc to other FEs
     public void pushVersionAsync(long dbId, OlapTable table, long version) {
         if (!Config.cloud_enable_version_syncer) {
@@ -153,17 +204,25 @@ public class CloudFEVersionSynchronizer {
         }
         // only update table version
         if (request.getPartitionVersionInfos().isEmpty()) {
+            List<OlapTable> invalidatedTables = new ArrayList<>();
             request.getTableVersionInfos().forEach(tableVersionInfo -> {
                 Table table = 
db.getTableNullable(tableVersionInfo.getTableId());
                 if (table == null || !table.isManagedTable()) {
                     return;
                 }
                 OlapTable olapTable = (OlapTable) table;
+                if (tableVersionInfo.getVersion() == INVALIDATE_VERSION_CACHE) 
{
+                    invalidatedTables.add(olapTable);
+                    return;
+                }
                 olapTable.setCachedTableVersion(tableVersionInfo.getVersion());
                 if (LOG.isDebugEnabled()) {
                     LOG.debug("Update tableId: {}, version: {}", 
olapTable.getId(), tableVersionInfo.getVersion());
                 }
             });
+            if (!invalidatedTables.isEmpty()) {
+                invalidateVersionCaches(invalidatedTables);
+            }
             return;
         }
         // update partition and table version
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudPartition.java 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudPartition.java
index cfb9fabe907..eb1193602ad 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudPartition.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudPartition.java
@@ -48,6 +48,7 @@ import java.util.Comparator;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.concurrent.atomic.AtomicLong;
 import java.util.concurrent.locks.ReentrantLock;
 import java.util.stream.Collectors;
 
@@ -66,6 +67,10 @@ public class CloudPartition extends Partition {
     // This value is set when get the version from meta-service, 0 means 
version is not cached yet
     private volatile long lastVersionCachedTimeMs = 0;
 
+    // Only an MS read started after invalidation may make this cache valid 
again.
+    private final AtomicLong versionCacheEpoch = new AtomicLong();
+    private final AtomicLong refreshedVersionCacheEpoch = new AtomicLong();
+
     private ReentrantLock lock = new ReentrantLock(true);
 
     public CloudPartition(long id, String name, MaterializedIndex baseIndex,
@@ -123,9 +128,13 @@ public class CloudPartition extends Partition {
         return super.getVisibleVersion();
     }
 
+    public void invalidateCachedVisibleVersion() {
+        versionCacheEpoch.incrementAndGet();
+    }
+
     @VisibleForTesting
     protected boolean isCachedVersionExpired() {
-        if (lastVersionCachedTimeMs == 0) {
+        if (lastVersionCachedTimeMs == 0 || versionCacheEpoch.get() != 
refreshedVersionCacheEpoch.get()) {
             return true;
         }
         ConnectContext ctx = ConnectContext.get();
@@ -151,6 +160,7 @@ public class CloudPartition extends Partition {
     }
 
     private long getVisibleVersionFromMs(boolean waitForPendingTxns) {
+        long cacheEpoch = versionCacheEpoch.get();
         if (LOG.isDebugEnabled()) {
             LOG.debug("getVisibleVersionFromMs use CloudPartition {}, 
waitForPendingTxns: {}",
                     super.getName(), waitForPendingTxns);
@@ -182,6 +192,7 @@ public class CloudPartition extends Partition {
                 LOG.debug("get version from meta service, version: {}, 
partition: {}", version, super.getId());
             }
             setCachedVisibleVersion(version, mTime);
+            refreshedVersionCacheEpoch.accumulateAndGet(cacheEpoch, Math::max);
             return version;
         } catch (RpcException e) {
             throw new RuntimeException("get version from meta service failed");
@@ -242,10 +253,12 @@ public class CloudPartition extends Partition {
         List<Long> tableIds = new ArrayList<>();
         List<Long> partitionIds = new ArrayList<>();
         List<Long> versionUpdateTimesMs = new ArrayList<>();
+        List<Long> cacheEpochs = new ArrayList<>();
         for (CloudPartition partition : partitions) {
             dbIds.add(partition.getDbId());
             tableIds.add(partition.getTableId());
             partitionIds.add(partition.getId());
+            cacheEpochs.add(partition.versionCacheEpoch.get());
         }
 
         List<Long> versions = getSnapshotVisibleVersion(
@@ -253,14 +266,26 @@ public class CloudPartition extends Partition {
 
         // Cache visible version, see hasData() for details.
         int size = versions.size();
-        for (int i = 0; i < size; ++i) {
-            Long version = versions.get(i);
-            if (version > Partition.PARTITION_INIT_VERSION) {
-                // For compatibility, the existing partitions may not have 
mtime
-                long mTime = versions.size() == versionUpdateTimesMs.size() ? 
versionUpdateTimesMs.get(i) : 0;
-                partitions.get(i).setCachedVisibleVersion(versions.get(i), 
mTime);
-            } else { // No data has been written to this partition
-                
partitions.get(i).setCachedVisibleVersion(Partition.PARTITION_INIT_VERSION, 
System.currentTimeMillis());
+        List<OlapTable> tables = getTables(partitions);
+        for (OlapTable table : tables) {
+            table.versionWriteLock();
+        }
+        try {
+            for (int i = 0; i < size; ++i) {
+                Long version = versions.get(i);
+                if (version > Partition.PARTITION_INIT_VERSION) {
+                    // For compatibility, the existing partitions may not have 
mtime
+                    long mTime = versions.size() == 
versionUpdateTimesMs.size() ? versionUpdateTimesMs.get(i) : 0;
+                    partitions.get(i).setCachedVisibleVersion(versions.get(i), 
mTime);
+                } else { // No data has been written to this partition
+                    
partitions.get(i).setCachedVisibleVersion(Partition.PARTITION_INIT_VERSION,
+                            System.currentTimeMillis());
+                }
+                
partitions.get(i).refreshedVersionCacheEpoch.accumulateAndGet(cacheEpochs.get(i),
 Math::max);
+            }
+        } finally {
+            for (int i = tables.size() - 1; i >= 0; i--) {
+                tables.get(i).versionWriteUnlock();
             }
         }
 
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgr.java
 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgr.java
index 1c49b3c11b2..9f8a133491e 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgr.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgr.java
@@ -478,6 +478,10 @@ public class CloudGlobalTransactionMgr implements 
GlobalTransactionMgrIface {
                         + "] is already aborted. abort reason: " + 
transactionState.getReason());
             } else if (transactionState.getTransactionStatus() == 
TransactionStatus.COMMITTED
                     || transactionState.getTransactionStatus() == 
TransactionStatus.VISIBLE) {
+                if (transactionState.getTransactionStatus() == 
TransactionStatus.VISIBLE) {
+                    ((CloudEnv) 
Env.getCurrentEnv()).getCloudFEVersionSynchronizer()
+                            .invalidateVersionCaches(dbId, 
transactionState.getTableIdList());
+                }
                 LOG.info("txn={}, status={} not need to calculate delete 
bitmap again, just return ",
                         transactionId,
                         transactionState.getTransactionStatus().toString());
@@ -582,6 +586,13 @@ public class CloudGlobalTransactionMgr implements 
GlobalTransactionMgrIface {
         long txnId = commitTxnResponse.getTxnInfo().getTxnId();
         int totalPartitionNum = commitTxnResponse.getPartitionIdsList().size();
         if (totalPartitionNum == 0 && 
commitTxnResponse.getTableStatsList().isEmpty()) {
+            // A visible retry can contain only txn_info. Invalidate locally 
without another commit-time RPC;
+            // subsequent reads fetch fresh versions, including the table 
version used to validate SQL caches.
+            if (commitTxnResponse.getTxnInfo().getStatus() == 
TxnStatusPB.TXN_STATUS_VISIBLE
+                    && !(commitTxnResponse.getIsLazyCommit() && 
commitTxnResponse.getIsLazyCommitIncomplete())) {
+                ((CloudEnv) 
Env.getCurrentEnv()).getCloudFEVersionSynchronizer()
+                        .invalidateVersionCaches(dbId, 
commitTxnResponse.getTxnInfo().getTableIdsList());
+            }
             return Collections.emptyMap();
         }
         Env env = Env.getCurrentEnv();
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgrTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgrTest.java
index 1daaaf51de8..490851d52db 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgrTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgrTest.java
@@ -21,7 +21,15 @@ import org.apache.doris.catalog.CatalogTestUtil;
 import org.apache.doris.catalog.Env;
 import org.apache.doris.catalog.FakeEditLog;
 import org.apache.doris.catalog.FakeEnv;
+import org.apache.doris.catalog.KeysType;
+import org.apache.doris.catalog.MaterializedIndex;
+import org.apache.doris.catalog.OlapTable;
+import org.apache.doris.catalog.PartitionInfo;
+import org.apache.doris.catalog.RandomDistributionInfo;
 import org.apache.doris.catalog.Table;
+import org.apache.doris.cloud.catalog.CloudEnv;
+import org.apache.doris.cloud.catalog.CloudFEVersionSynchronizer;
+import org.apache.doris.cloud.catalog.CloudPartition;
 import org.apache.doris.cloud.proto.Cloud;
 import org.apache.doris.cloud.proto.Cloud.AbortTxnResponse;
 import org.apache.doris.cloud.proto.Cloud.BeginTxnResponse;
@@ -31,29 +39,47 @@ import 
org.apache.doris.cloud.proto.Cloud.GetCurrentMaxTxnResponse;
 import org.apache.doris.cloud.proto.Cloud.MetaServiceCode;
 import org.apache.doris.cloud.proto.Cloud.TxnInfoPB;
 import org.apache.doris.cloud.rpc.MetaServiceProxy;
+import org.apache.doris.cloud.rpc.VersionHelper;
 import org.apache.doris.common.AnalysisException;
+import org.apache.doris.common.ClientPool;
 import org.apache.doris.common.Config;
 import org.apache.doris.common.DuplicatedRequestException;
 import org.apache.doris.common.FeMetaVersion;
+import org.apache.doris.common.GenericPool;
 import org.apache.doris.common.LabelAlreadyUsedException;
 import org.apache.doris.common.MetaNotFoundException;
 import org.apache.doris.common.QuotaExceedException;
 import org.apache.doris.common.UserException;
+import org.apache.doris.common.jmockit.Deencapsulation;
 import org.apache.doris.load.routineload.RLTaskTxnCommitAttachment;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.SessionVariable;
+import org.apache.doris.rpc.RpcException;
+import org.apache.doris.service.FrontendServiceImpl;
+import org.apache.doris.system.Frontend;
+import org.apache.doris.system.SystemInfoService.HostInfo;
+import org.apache.doris.thrift.FrontendService;
+import org.apache.doris.thrift.TCloudVersionInfo;
+import org.apache.doris.thrift.TFrontendSyncCloudVersionRequest;
+import org.apache.doris.thrift.TStatus;
+import org.apache.doris.thrift.TStatusCode;
 import org.apache.doris.thrift.TTabletCommitInfo;
 import org.apache.doris.transaction.BeginTransactionException;
 import org.apache.doris.transaction.GlobalTransactionMgrIface;
 import org.apache.doris.transaction.TabletCommitInfo;
 import org.apache.doris.transaction.TransactionState;
+import org.apache.doris.transaction.TransactionStatus;
 import org.apache.doris.transaction.TxnStateChangeCallback;
 
 import com.google.common.collect.Lists;
 import mockit.Mock;
 import mockit.MockUp;
+import org.junit.After;
 import org.junit.Assert;
 import org.junit.Before;
 import org.junit.Test;
 import org.junit.jupiter.api.Assertions;
+import org.mockito.AdditionalAnswers;
 import org.mockito.ArgumentCaptor;
 import org.mockito.MockedStatic;
 import org.mockito.Mockito;
@@ -62,8 +88,15 @@ import java.lang.reflect.InvocationTargetException;
 import java.lang.reflect.Method;
 import java.util.List;
 import java.util.Map;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.atomic.AtomicReference;
 
 public class CloudGlobalTransactionMgrTest {
 
@@ -83,11 +116,21 @@ public class CloudGlobalTransactionMgrTest {
         Config.meta_service_endpoint = "127.0.0.1:20121";
         fakeEditLog = new FakeEditLog();
         fakeEnv = new FakeEnv();
-        masterEnv = CatalogTestUtil.createTestCatalog();
+        Env catalog = CatalogTestUtil.createTestCatalog();
+        masterEnv = Mockito.mock(CloudEnv.class, 
AdditionalAnswers.delegatesTo(catalog));
+        Mockito.doReturn(new CloudFEVersionSynchronizer()).when((CloudEnv) 
masterEnv).getCloudFEVersionSynchronizer();
+        // Env.getCurrentGlobalTransactionMgr reads the field directly rather 
than calling the delegated getter.
+        Deencapsulation.setField(masterEnv, "globalTransactionMgr", 
catalog.getGlobalTransactionMgr());
+        FakeEnv.setEnv(masterEnv);
         FakeEnv.setMetaVersion(FeMetaVersion.VERSION_CURRENT);
         masterTransMgr = masterEnv.getGlobalTransactionMgr();
     }
 
+    @After
+    public void tearDown() {
+        ConnectContext.remove();
+    }
+
     @Test
     public void testBeginTransaction() throws LabelAlreadyUsedException, 
AnalysisException,
             BeginTransactionException, DuplicatedRequestException, 
QuotaExceedException, MetaNotFoundException {
@@ -554,6 +597,361 @@ public class CloudGlobalTransactionMgrTest {
         Assert.assertEquals(1000, result);
     }
 
+    @Test
+    public void testVisibleRetryInvalidatesTableAndPartitionVersions() throws 
Exception {
+        useVersionCaches();
+        CloudPartition first = addCloudPartition(1000);
+        CloudPartition second = addCloudPartition(2000);
+        OlapTable firstTable = getCloudTable(first);
+        OlapTable secondTable = getCloudTable(second);
+        CommitTxnResponse response = 
visibleRetry(Lists.newArrayList(first.getTableId(), second.getTableId()));
+        try (MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+            versions.when(() -> 
VersionHelper.getVersionFromMeta(Mockito.any())).thenReturn(partitionVersion(4));
+            masterTransMgr.afterCommitTxnResp(response, null, 
Lists.newArrayList());
+            versions.verifyNoInteractions();
+            Assertions.assertEquals(2, first.getCachedVisibleVersion());
+            Assertions.assertEquals(2, firstTable.getCachedTableVersion());
+            Assertions.assertEquals(4, first.getVisibleVersion());
+            Assertions.assertEquals(Lists.newArrayList(4L), 
CloudPartition.getSnapshotVisibleVersion(Lists.newArrayList(second)));
+            Assertions.assertEquals(4, firstTable.getVisibleVersion());
+            Assertions.assertEquals(Lists.newArrayList(4L), 
OlapTable.getVisibleVersionInBatch(Lists.newArrayList(secondTable)));
+            Assertions.assertEquals(4, first.getCachedVisibleVersion());
+            Assertions.assertEquals(4, second.getCachedVisibleVersion());
+            Assertions.assertEquals(40, first.getVisibleVersionTime());
+            // Current versions must not become the original transaction's BE 
promotion outcome.
+            Assertions.assertEquals(0, response.getVersionsCount());
+            Assertions.assertEquals(4, first.getVisibleVersion());
+            Assertions.assertEquals(Lists.newArrayList(4L), 
CloudPartition.getSnapshotVisibleVersion(Lists.newArrayList(second)));
+            Assertions.assertEquals(4, firstTable.getVisibleVersion());
+            Assertions.assertEquals(Lists.newArrayList(4L), 
OlapTable.getVisibleVersionInBatch(Lists.newArrayList(secondTable)));
+            versions.verify(() -> 
VersionHelper.getVersionFromMeta(Mockito.any()), Mockito.times(4));
+        }
+    }
+
+    @Test
+    public void testIncompleteCommitDoesNotRefreshPartitionVersions() throws 
Exception {
+        CloudPartition partition = addCloudPartition(1000);
+        CommitTxnResponse response = 
visibleRetry(Lists.newArrayList(partition.getTableId()));
+        try (MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+            
masterTransMgr.afterCommitTxnResp(response.toBuilder().setTxnInfo(response.getTxnInfo().toBuilder()
+                    
.setStatus(Cloud.TxnStatusPB.TXN_STATUS_COMMITTED)).build(), null, 
Lists.newArrayList());
+            
masterTransMgr.afterCommitTxnResp(response.toBuilder().setIsLazyCommit(true)
+                    .setIsLazyCommitIncomplete(true).build(), null, 
Lists.newArrayList());
+            Assertions.assertEquals(2, partition.getCachedVisibleVersion());
+            versions.verifyNoInteractions();
+        }
+    }
+
+    @Test
+    public void testVisibleRetrySkipsDroppedTable() throws Exception {
+        try (MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+            
masterTransMgr.afterCommitTxnResp(visibleRetry(Lists.newArrayList(1000L)), 
null, Lists.newArrayList());
+            versions.verifyNoInteractions();
+        }
+    }
+
+    @Test
+    public void testMowVisiblePrecheckInvalidatesVersions() throws Exception {
+        useVersionCaches();
+        CloudPartition partition = addCloudPartition(1000);
+        TxnInfoPB txnInfo = 
visibleRetry(Lists.newArrayList(partition.getTableId())).getTxnInfo();
+        MetaServiceProxy proxy = Mockito.mock(MetaServiceProxy.class);
+        
Mockito.when(proxy.getTxn(Mockito.any())).thenReturn(Cloud.GetTxnResponse.newBuilder()
+                
.setStatus(Cloud.MetaServiceResponseStatus.newBuilder().setCode(MetaServiceCode.OK))
+                .setTxnInfo(txnInfo).build());
+        try (MockedStatic<MetaServiceProxy> proxyMock = 
Mockito.mockStatic(MetaServiceProxy.class);
+                MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+            proxyMock.when(MetaServiceProxy::getInstance).thenReturn(proxy);
+            versions.when(() -> 
VersionHelper.getVersionFromMeta(Mockito.any())).thenReturn(partitionVersion(4));
+            Method check = 
CloudGlobalTransactionMgr.class.getDeclaredMethod("checkTransactionStateBeforeCommit",
+                    long.class, long.class);
+            check.setAccessible(true);
+            Assertions.assertEquals(false, check.invoke(masterTransMgr, 
txnInfo.getDbId(), txnInfo.getTxnId()));
+            versions.verifyNoInteractions();
+            Assertions.assertEquals(4, partition.getVisibleVersion());
+            Assertions.assertEquals(4, 
getCloudTable(partition).getVisibleVersion());
+            Assertions.assertEquals(4, partition.getCachedVisibleVersion());
+        }
+    }
+
+    @Test
+    public void testRefreshFailurePreservesCommitCallbacks() throws Exception {
+        useVersionCaches();
+        CloudPartition partition = addCloudPartition(1000);
+        CommitTxnResponse retry = 
visibleRetry(Lists.newArrayList(partition.getTableId()));
+        CommitTxnResponse response = retry.toBuilder()
+                
.setTxnInfo(retry.getTxnInfo().toBuilder().setListenerId(42)).build();
+        TxnStateChangeCallback callback = 
Mockito.mock(TxnStateChangeCallback.class);
+        Mockito.when(callback.getId()).thenReturn(42L);
+        masterTransMgr.getCallbackFactory().addCallback(callback);
+        Mockito.clearInvocations(callback);
+        MetaServiceProxy proxy = Mockito.mock(MetaServiceProxy.class);
+        Mockito.when(proxy.commitTxn(Mockito.any())).thenReturn(response);
+        Table table = 
masterEnv.getInternalCatalog().getDbOrMetaException(CatalogTestUtil.testDbId1)
+                .getTableOrMetaException(partition.getTableId());
+        try (MockedStatic<MetaServiceProxy> proxyMock = 
Mockito.mockStatic(MetaServiceProxy.class);
+                MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+            proxyMock.when(MetaServiceProxy::getInstance).thenReturn(proxy);
+            versions.when(() -> 
VersionHelper.getVersionFromMeta(Mockito.any()))
+                    .thenThrow(new RpcException("MS", "unavailable"));
+            Assertions.assertDoesNotThrow(() -> 
masterTransMgr.commitTransactionWithoutLock(CatalogTestUtil.testDbId1,
+                    Lists.newArrayList(table), 
response.getTxnInfo().getTxnId(), null, null));
+            Assertions.assertEquals(2, partition.getCachedVisibleVersion());
+            versions.verifyNoInteractions();
+            Mockito.verify(callback).afterCommitted(Mockito.argThat(state ->
+                    state.getTransactionStatus() == 
TransactionStatus.VISIBLE), Mockito.eq(true));
+            Mockito.verify(callback).afterVisible(Mockito.argThat(state ->
+                    state.getTransactionStatus() == 
TransactionStatus.VISIBLE), Mockito.eq(true));
+            // A version outage fails the following read, not the already 
committed transaction.
+            Assertions.assertThrows(RuntimeException.class, 
partition::getVisibleVersion);
+            Assertions.assertThrows(RpcException.class, ((OlapTable) 
table)::getVisibleVersion);
+            versions.when(() -> 
VersionHelper.getVersionFromMeta(Mockito.any())).thenReturn(partitionVersion(4));
+            Assertions.assertEquals(4, partition.getVisibleVersion());
+            Assertions.assertEquals(4, ((OlapTable) 
table).getVisibleVersion());
+            Mockito.verifyNoMoreInteractions(callback);
+        }
+    }
+
+    @Test
+    public void testSinglePartitionReadCannotClearConcurrentInvalidation() 
throws Exception {
+        checkReadCannotClearConcurrentInvalidation(false, false);
+    }
+
+    @Test
+    public void testBatchPartitionReadCannotClearConcurrentInvalidation() 
throws Exception {
+        checkReadCannotClearConcurrentInvalidation(false, true);
+    }
+
+    @Test
+    public void testSingleTableReadCannotClearConcurrentInvalidation() throws 
Exception {
+        checkReadCannotClearConcurrentInvalidation(true, false);
+    }
+
+    @Test
+    public void testBatchTableReadCannotClearConcurrentInvalidation() throws 
Exception {
+        checkReadCannotClearConcurrentInvalidation(true, true);
+    }
+
+    private void checkReadCannotClearConcurrentInvalidation(boolean 
tableVersion, boolean batch) throws Exception {
+        useVersionCaches();
+        CloudPartition partition = addCloudPartition(1000);
+        OlapTable table = getCloudTable(partition);
+        
ConnectContext.get().getSessionVariable().cloudPartitionVersionCacheTtlMs = 0;
+        ConnectContext.get().getSessionVariable().cloudTableVersionCacheTtlMs 
= 0;
+        try (MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+            versions.when(() -> 
VersionHelper.getVersionFromMeta(Mockito.any())).thenAnswer(invocation -> {
+                // The MS snapshot predates the commit, but its reply arrives 
after cache invalidation.
+                
masterTransMgr.afterCommitTxnResp(visibleRetry(Lists.newArrayList(table.getId())),
 null, Lists.newArrayList());
+                useVersionCaches();
+                return partitionVersion(2);
+            });
+            if (tableVersion) {
+                Assertions.assertEquals(2L, batch
+                        ? 
OlapTable.getVisibleVersionInBatch(Lists.newArrayList(table)).get(0) : 
table.getVisibleVersion());
+            } else {
+                Assertions.assertEquals(2L, batch
+                        ? 
CloudPartition.getSnapshotVisibleVersion(Lists.newArrayList(partition)).get(0)
+                        : partition.getVisibleVersion());
+            }
+            // A delayed ordinary commit notification cannot repair an unknown 
version either.
+            table.setCachedTableVersion(3);
+            partition.setCachedVisibleVersion(3, 30);
+            versions.when(() -> 
VersionHelper.getVersionFromMeta(Mockito.any())).thenReturn(partitionVersion(4));
+            if (tableVersion) {
+                Assertions.assertEquals(4, table.getVisibleVersion());
+                Assertions.assertEquals(4, table.getVisibleVersion());
+            } else {
+                Assertions.assertEquals(4, partition.getVisibleVersion());
+                Assertions.assertEquals(4, partition.getVisibleVersion());
+            }
+            versions.verify(() -> 
VersionHelper.getVersionFromMeta(Mockito.any()), Mockito.times(2));
+        }
+    }
+
+    @Test
+    public void testInvalidationWaitsForVersionSnapshotReaders() throws 
Exception {
+        CloudPartition partition = addCloudPartition(1000);
+        OlapTable table = getCloudTable(partition);
+        ExecutorService executor = Executors.newSingleThreadExecutor();
+        CountDownLatch started = new CountDownLatch(1);
+        table.versionReadLock();
+        CompletableFuture<Void> invalidation;
+        try {
+            invalidation = CompletableFuture.runAsync(() -> {
+                started.countDown();
+                ((CloudEnv) masterEnv).getCloudFEVersionSynchronizer()
+                        .invalidateVersionCaches(CatalogTestUtil.testDbId1, 
Lists.newArrayList(table.getId()));
+            }, executor);
+            Assertions.assertTrue(started.await(5, TimeUnit.SECONDS));
+            Assertions.assertThrows(TimeoutException.class, () -> 
invalidation.get(200, TimeUnit.MILLISECONDS));
+        } finally {
+            table.versionReadUnlock();
+            executor.shutdown();
+        }
+        invalidation.get(5, TimeUnit.SECONDS);
+        Assertions.assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS));
+    }
+
+    @Test
+    public void testBatchRefreshInstallsVersionsUnderWriteLock() throws 
Exception {
+        CloudPartition first = addCloudPartition(1000);
+        OlapTable table = getCloudTable(first);
+        CloudPartition second = Mockito.spy(new CloudPartition(1002, "p2", new 
MaterializedIndex(),
+                new RandomDistributionInfo(1), CatalogTestUtil.testDbId1, 
table.getId()));
+        table.addPartition(second);
+        ExecutorService executor = Executors.newSingleThreadExecutor();
+        AtomicReference<CompletableFuture<Void>> snapshotReader = new 
AtomicReference<>();
+        try (MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+            versions.when(() -> 
VersionHelper.getVersionFromMeta(Mockito.any())).thenReturn(
+                    
partitionVersion(4).toBuilder().addVersions(4).addVersionUpdateTimeMs(40).build());
+            Mockito.doAnswer(invocation -> {
+                CountDownLatch started = new CountDownLatch(1);
+                CompletableFuture<Void> reader = CompletableFuture.runAsync(() 
-> {
+                    started.countDown();
+                    table.versionReadLock();
+                    try {
+                        Assertions.assertEquals(4, 
first.getCachedVisibleVersion());
+                        Assertions.assertEquals(4, 
second.getCachedVisibleVersion());
+                    } finally {
+                        table.versionReadUnlock();
+                    }
+                }, executor);
+                snapshotReader.set(reader);
+                Assertions.assertTrue(started.await(5, TimeUnit.SECONDS));
+                Assertions.assertThrows(TimeoutException.class, () -> 
reader.get(200, TimeUnit.MILLISECONDS));
+                return invocation.callRealMethod();
+            }).when(second).setCachedVisibleVersion(4, 40);
+            
CloudPartition.getSnapshotVisibleVersionFromMs(Lists.newArrayList(first, 
second), false);
+            snapshotReader.get().get(5, TimeUnit.SECONDS);
+        } finally {
+            executor.shutdown();
+            Assertions.assertTrue(executor.awaitTermination(5, 
TimeUnit.SECONDS));
+        }
+    }
+
+    @Test
+    public void testFollowerInvalidatesTableAndAllPartitions() throws 
Exception {
+        useVersionCaches();
+        CloudPartition partition = addCloudPartition(1000);
+        OlapTable table = getCloudTable(partition);
+        CloudPartition second = new CloudPartition(1002, "p2", new 
MaterializedIndex(),
+                new RandomDistributionInfo(1), CatalogTestUtil.testDbId1, 
table.getId());
+        second.setCachedVisibleVersion(2, 20);
+        table.addPartition(second);
+        TFrontendSyncCloudVersionRequest request = new 
TFrontendSyncCloudVersionRequest()
+                
.setDbId(CatalogTestUtil.testDbId1).setPartitionVersionInfos(Lists.newArrayList())
+                .setTableVersionInfos(Lists.newArrayList(new 
TCloudVersionInfo().setTableId(table.getId()).setVersion(-1)));
+        Mockito.doReturn(false).when(masterEnv).isMaster();
+        FrontendServiceImpl service = new FrontendServiceImpl(null);
+        try (MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+            Assertions.assertEquals(TStatusCode.OK, 
service.syncCloudVersion(request).getStatusCode());
+            versions.verifyNoInteractions();
+            // Even a delayed regular version push must leave these caches 
invalid.
+            service.syncCloudVersion(new TFrontendSyncCloudVersionRequest()
+                    
.setDbId(CatalogTestUtil.testDbId1).setPartitionVersionInfos(Lists.newArrayList())
+                    .setTableVersionInfos(Lists.newArrayList(new 
TCloudVersionInfo().setTableId(table.getId()).setVersion(3))));
+            versions.when(() -> 
VersionHelper.getVersionFromMeta(Mockito.any())).thenReturn(partitionVersion(4));
+            Assertions.assertEquals(4, table.getVisibleVersion());
+            Assertions.assertEquals(4, partition.getVisibleVersion());
+            Assertions.assertEquals(4, second.getVisibleVersion());
+            versions.verify(() -> 
VersionHelper.getVersionFromMeta(Mockito.any()), Mockito.times(3));
+        }
+    }
+
+    @Test
+    @SuppressWarnings("unchecked")
+    public void testInvalidationIsPushedThroughVersionRpc() throws Exception {
+        CloudPartition partition = addCloudPartition(1000);
+        boolean syncEnabled = Config.cloud_enable_version_syncer;
+        GenericPool<FrontendService.Client> originalPool = 
ClientPool.frontendVersionPool;
+        GenericPool<FrontendService.Client> pool = 
Mockito.mock(GenericPool.class);
+        FrontendService.Client client = 
Mockito.mock(FrontendService.Client.class);
+        Mockito.when(pool.borrowObject(Mockito.any())).thenReturn(client);
+        CompletableFuture<TFrontendSyncCloudVersionRequest> sent = new 
CompletableFuture<>();
+        CountDownLatch returned = new CountDownLatch(1);
+        
Mockito.when(client.syncCloudVersion(Mockito.any())).thenAnswer(invocation -> {
+            sent.complete(invocation.getArgument(0));
+            return new TStatus(TStatusCode.OK);
+        });
+        Mockito.doAnswer(invocation -> {
+            returned.countDown();
+            return null;
+        }).when(pool).returnObject(Mockito.any(), Mockito.eq(client));
+        CloudEnv sender = (CloudEnv) masterEnv;
+        Frontend follower = Mockito.mock(Frontend.class);
+        Mockito.when(follower.isAlive()).thenReturn(true);
+        Mockito.when(follower.getHost()).thenReturn("127.0.0.2");
+        Mockito.when(follower.getRpcPort()).thenReturn(9020);
+        
Mockito.doReturn(Lists.newArrayList(follower)).when(sender).getFrontends(null);
+        Mockito.doReturn(new HostInfo("127.0.0.1", 
9010)).when(sender).getSelfNode();
+        try {
+            FakeEnv.setEnv(sender);
+            Config.cloud_enable_version_syncer = true;
+            ClientPool.frontendVersionPool = pool;
+            
masterTransMgr.afterCommitTxnResp(visibleRetry(Lists.newArrayList(partition.getTableId())),
 null, Lists.newArrayList());
+            TFrontendSyncCloudVersionRequest request = sent.get(5, 
TimeUnit.SECONDS);
+            Assertions.assertTrue(returned.await(5, TimeUnit.SECONDS));
+            // Check the actual Thrift payload as well as the sender's Java 
objects.
+            TFrontendSyncCloudVersionRequest decoded = new 
TFrontendSyncCloudVersionRequest();
+            new org.apache.thrift.TDeserializer().deserialize(decoded,
+                    new org.apache.thrift.TSerializer().serialize(request));
+            Assertions.assertEquals(CatalogTestUtil.testDbId1, 
decoded.getDbId());
+            
Assertions.assertTrue(decoded.getPartitionVersionInfos().isEmpty());
+            Assertions.assertEquals(1, decoded.getTableVersionInfosSize());
+            Assertions.assertEquals(partition.getTableId(), 
decoded.getTableVersionInfos().get(0).getTableId());
+            Assertions.assertEquals(-1, 
decoded.getTableVersionInfos().get(0).getVersion());
+        } finally {
+            FakeEnv.setEnv(masterEnv);
+            ClientPool.frontendVersionPool = originalPool;
+            Config.cloud_enable_version_syncer = syncEnabled;
+        }
+    }
+
+    private CloudPartition addCloudPartition(long tableId) throws Exception {
+        OlapTable table = new OlapTable(tableId, "version_cache_" + tableId, 
Lists.newArrayList(), KeysType.DUP_KEYS,
+                new PartitionInfo(), new RandomDistributionInfo(1));
+        CloudPartition partition = new CloudPartition(tableId + 1, "p1", new 
MaterializedIndex(),
+                new RandomDistributionInfo(1), CatalogTestUtil.testDbId1, 
tableId);
+        partition.setCachedVisibleVersion(2, 20);
+        table.addPartition(partition);
+        table.setCachedTableVersion(2);
+        
masterEnv.getInternalCatalog().getDbOrMetaException(CatalogTestUtil.testDbId1).registerTable(table);
+        return partition;
+    }
+
+    private MockedStatic<VersionHelper> mockVersionHelper() {
+        MockedStatic<VersionHelper> versions = 
Mockito.mockStatic(VersionHelper.class);
+        // Batch reads pass an explicit retry limit; share the response stub 
with single-version reads.
+        versions.when(() -> VersionHelper.getVersionFromMeta(Mockito.any(), 
Mockito.anyInt()))
+                .thenAnswer(invocation -> 
VersionHelper.getVersionFromMeta(invocation.getArgument(0)));
+        return versions;
+    }
+
+    private CommitTxnResponse visibleRetry(List<Long> tableIds) {
+        return CommitTxnResponse.newBuilder()
+                
.setStatus(Cloud.MetaServiceResponseStatus.newBuilder().setCode(MetaServiceCode.OK))
+                
.setTxnInfo(buildTxnInfo(100).toBuilder().clearTableIds().addAllTableIds(tableIds)
+                        
.setStatus(Cloud.TxnStatusPB.TXN_STATUS_VISIBLE)).build();
+    }
+
+    private Cloud.GetVersionResponse partitionVersion(long version) {
+        return 
Cloud.GetVersionResponse.newBuilder().setVersion(version).addVersions(version)
+                .addVersionUpdateTimeMs(version * 10).build();
+    }
+
+    private OlapTable getCloudTable(CloudPartition partition) throws Exception 
{
+        return (OlapTable) 
masterEnv.getInternalCatalog().getDbOrMetaException(CatalogTestUtil.testDbId1)
+                .getTableOrMetaException(partition.getTableId());
+    }
+
+    private void useVersionCaches() {
+        ConnectContext context = new ConnectContext();
+        context.setSessionVariable(new SessionVariable());
+        context.getSessionVariable().cloudPartitionVersionCacheTtlMs = 
Long.MAX_VALUE;
+        context.getSessionVariable().cloudTableVersionCacheTtlMs = 
Long.MAX_VALUE;
+        context.setThreadLocalInfo();
+    }
+
     private TxnInfoPB buildTxnInfo(long transactionId) {
         return TxnInfoPB.newBuilder()
                 .setDbId(CatalogTestUtil.testDbId1)


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to