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]