This is an automated email from the ASF dual-hosted git repository.
luwei16 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new 3dbe4d3ca53 [fix](cloud) Invalidate version caches on visible commit
retries (#67813)
3dbe4d3ca53 is described below
commit 3dbe4d3ca53307f73db96eda23d4efc9f02963c5
Author: Luwei <[email protected]>
AuthorDate: Wed Sep 16 12:42:59 2026 +0800
[fix](cloud) Invalidate version caches on visible commit retries (#67813)
### What problem does this PR solve?
Issue Number: close #67099
Related PR: None
Problem Summary: After a cloud commit response is lost, a retry can
return a VISIBLE transaction without version results. FE can then keep
using stale partition versions and the table version used to validate
SQL caches.
Invalidate the affected table and partition version caches, including
the MoW already-visible shortcut, and notify other FEs through the
existing version synchronization RPC. Subsequent reads fetch fresh
versions from MS. Commit processing performs no extra MS version query,
so a version-service outage does not turn an already committed
transaction into a commit error or require replaying its callbacks.
Publish batch cache updates under the tables' version write locks.
Invalidation epochs prevent old in-flight MS responses and delayed
commit notifications from restoring stale caches. The synchronization
RPC uses table version -1 for invalidation; older FEs ignore it and
retain their existing periodic synchronization. Peer notification keeps
its existing asynchronous behavior and configuration.
### Release note
Prevent stale reads and false commit failures after cloud commit
retries.
### Check List (For Author)
- Test: Unit Test
- 58 FE tests passed in both ordinary and JaCoCo coverage modes via
run-fe-ut.sh (CloudGlobalTransactionMgrTest, CloudPartitionTest,
OlapTableTest, VersionHelperTest); FE Checkstyle: 0 violations; PR diff
whitespace check passed.
- Tests cover invalidation and batch-publication locks, table/partition
cache recovery, old RPC completions, successful commit callbacks during
a version outage, and FE notification serialization/receiver handling.
- No live cluster deployment or SQL regression run.
- Behavior changed: Yes. VISIBLE retries without versions invalidate
local and peer caches; subsequent reads retrieve versions without
changing the successful commit result.
- Does this need documentation: No
---
.../java/org/apache/doris/catalog/OlapTable.java | 20 +-
.../cloud/catalog/CloudFEVersionSynchronizer.java | 59 +++
.../apache/doris/cloud/catalog/CloudPartition.java | 46 ++-
.../transaction/CloudGlobalTransactionMgr.java | 11 +
.../transaction/CloudGlobalTransactionMgrTest.java | 398 ++++++++++++++++++++-
5 files changed, 520 insertions(+), 14 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 17cac68dc68..4645a7796c9 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
@@ -248,6 +248,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() {
@@ -3591,7 +3595,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();
@@ -3605,13 +3609,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();
@@ -3632,6 +3641,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())
@@ -3658,6 +3668,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);
@@ -3722,9 +3733,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);
@@ -3733,6 +3746,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 5b8f5507a34..6ac2abbbace 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) {
@@ -154,17 +205,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 86c3b72409c..07c2a6d504e 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,
@@ -127,9 +132,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();
@@ -155,6 +164,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);
@@ -183,6 +193,7 @@ public class CloudPartition extends Partition {
version, tso, super.getId());
}
setCachedVisibleVersion(version, mTime, tso);
+ refreshedVersionCacheEpoch.accumulateAndGet(cacheEpoch,
Math::max);
return version;
} else {
assert resp.getStatus().getCode() ==
MetaServiceCode.VERSION_NOT_FOUND;
@@ -193,6 +204,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");
@@ -254,10 +266,12 @@ public class CloudPartition extends Partition {
List<Long> partitionIds = new ArrayList<>();
List<Long> versionUpdateTimesMs = new ArrayList<>();
List<Long> commitTsos = 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(
@@ -266,15 +280,27 @@ public class CloudPartition extends Partition {
// Cache visible version, see hasData() for details.
int size = versions.size();
boolean hasCommitTsos = commitTsos.size() == 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;
- long tso = hasCommitTsos ? commitTsos.get(i) : -1;
- partitions.get(i).setCachedVisibleVersion(versions.get(i),
mTime, tso);
- } 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;
+ long tso = hasCommitTsos ? commitTsos.get(i) : -1;
+ partitions.get(i).setCachedVisibleVersion(versions.get(i),
mTime, tso);
+ } 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 52b7a84bea2..0aa654c80fc 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
@@ -491,6 +491,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());
@@ -597,6 +601,13 @@ public class CloudGlobalTransactionMgr implements
GlobalTransactionMgrIface {
? commitTxnResponse.getTxnInfo().getCommitTso() : -1;
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 3d200ce80d3..b2927e26eba 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,9 +21,17 @@ 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.catalog.stream.CloudOlapTableStreamUpdate;
import org.apache.doris.catalog.stream.TableStreamUpdateInfo;
+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;
@@ -33,16 +41,32 @@ 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.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.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.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;
@@ -50,6 +74,7 @@ import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
+import org.mockito.AdditionalAnswers;
import org.mockito.ArgumentCaptor;
import org.mockito.MockedStatic;
import org.mockito.Mockito;
@@ -57,8 +82,15 @@ import org.mockito.Mockito;
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 {
@@ -77,13 +109,19 @@ 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();
}
@AfterEach
public void tearDown() {
+ ConnectContext.remove();
if (fakeEnv != null) {
fakeEnv.close();
}
@@ -641,6 +679,364 @@ public class CloudGlobalTransactionMgrTest {
}
}
+ @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(List.of(first.getTableId(),
second.getTableId()));
+ try (MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+ versions.when(() ->
VersionHelper.getVersionFromMeta(Mockito.any())).thenReturn(partitionVersion(4));
+ masterTransMgr.afterCommitTxnResp(response, null, List.of());
+ versions.verifyNoInteractions();
+ Assertions.assertEquals(2, first.getCachedVisibleVersion());
+ Assertions.assertEquals(2, firstTable.getCachedTableVersion());
+ Assertions.assertEquals(4, first.getVisibleVersion());
+ Assertions.assertEquals(List.of(4L),
CloudPartition.getSnapshotVisibleVersion(List.of(second)));
+ Assertions.assertEquals(4, firstTable.getVisibleVersion());
+ Assertions.assertEquals(List.of(4L),
OlapTable.getVisibleVersionInBatch(List.of(secondTable)));
+ Assertions.assertEquals(4, first.getCachedVisibleVersion());
+ Assertions.assertEquals(4, second.getCachedVisibleVersion());
+ Assertions.assertEquals(40, first.getVisibleVersionTime());
+ Assertions.assertEquals(400, first.getTso());
+ // Current versions must not become the original transaction's BE
promotion outcome.
+ Assertions.assertEquals(0, response.getVersionsCount());
+ Assertions.assertEquals(4, first.getVisibleVersion());
+ Assertions.assertEquals(List.of(4L),
CloudPartition.getSnapshotVisibleVersion(List.of(second)));
+ Assertions.assertEquals(4, firstTable.getVisibleVersion());
+ Assertions.assertEquals(List.of(4L),
OlapTable.getVisibleVersionInBatch(List.of(secondTable)));
+ versions.verify(() ->
VersionHelper.getVersionFromMeta(Mockito.any()), Mockito.times(4));
+ }
+ }
+
+ @Test
+ public void testIncompleteCommitDoesNotRefreshPartitionVersions() throws
Exception {
+ CloudPartition partition = addCloudPartition(1000);
+ CommitTxnResponse response =
visibleRetry(List.of(partition.getTableId()));
+ try (MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+
masterTransMgr.afterCommitTxnResp(response.toBuilder().setTxnInfo(response.getTxnInfo().toBuilder()
+
.setStatus(Cloud.TxnStatusPB.TXN_STATUS_COMMITTED)).build(), null, List.of());
+
masterTransMgr.afterCommitTxnResp(response.toBuilder().setIsLazyCommit(true)
+ .setIsLazyCommitIncomplete(true).build(), null, List.of());
+ Assertions.assertEquals(2, partition.getCachedVisibleVersion());
+ versions.verifyNoInteractions();
+ }
+ }
+
+ @Test
+ public void testVisibleRetrySkipsDroppedTable() throws Exception {
+ try (MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+ masterTransMgr.afterCommitTxnResp(visibleRetry(List.of(1000L)),
null, List.of());
+ versions.verifyNoInteractions();
+ }
+ }
+
+ @Test
+ public void testMowVisiblePrecheckInvalidatesVersions() throws Exception {
+ useVersionCaches();
+ CloudPartition partition = addCloudPartition(1000);
+ TxnInfoPB txnInfo =
visibleRetry(List.of(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(List.of(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,
+ List.of(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(List.of(table.getId())), null,
List.of());
+ useVersionCaches();
+ return partitionVersion(2);
+ });
+ if (tableVersion) {
+ Assertions.assertEquals(2L, batch
+ ?
OlapTable.getVisibleVersionInBatch(List.of(table)).get(0) :
table.getVisibleVersion());
+ } else {
+ Assertions.assertEquals(2L, batch
+ ?
CloudPartition.getSnapshotVisibleVersion(List.of(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(() -> {
+ try (FakeEnv ignored = new FakeEnv()) {
+ started.countDown();
+ ((CloudEnv) masterEnv).getCloudFEVersionSynchronizer()
+
.invalidateVersionCaches(CatalogTestUtil.testDbId1, List.of(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).addCommitTsos(400).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, 400);
+ CloudPartition.getSnapshotVisibleVersionFromMs(List.of(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(List.of())
+ .setTableVersionInfos(List.of(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(List.of())
+ .setTableVersionInfos(List.of(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(List.of(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(List.of(partition.getTableId())),
null, List.of());
+ 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,
List.of(), 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, 200);
+ 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).setCommitTso(300)).build();
+ }
+
+ private Cloud.GetVersionResponse partitionVersion(long version) {
+ return
Cloud.GetVersionResponse.newBuilder().setVersion(version).addVersions(version)
+ .addVersionUpdateTimeMs(version * 10).addCommitTsos(version *
100).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]