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 e5a4e725fac [fix](cloud) Wait for running transactions before
incremental reads (#67181)
e5a4e725fac is described below
commit e5a4e725facc5ebc99b576ebda9ab178691dac19
Author: Luwei <[email protected]>
AuthorDate: Wed Sep 2 14:16:40 2026 +0800
[fix](cloud) Wait for running transactions before incremental reads (#67181)
### What problem does this PR solve?
Issue Number: None
Related PR: None
Problem Summary: In cloud mode, time-based incremental reads treated an
empty committed-transaction list as proof that a read window was
complete. A transaction could already have a commit timestamp in the
window while its delete bitmap and partition version were still being
published, allowing the query to return an empty result and downstream
consumers to close the window. Capture a MetaService transaction ID
watermark at query start, wait for earlier target-table transactions to
finish through the existing conflict check, and fetch fresh visible
versions for incremental scans after the wait.
### Release note
Cloud time-based incremental reads now wait for transactions registered
before query start to finish and read the latest visible partition
versions before scanning.
### Check List (For Author)
- Test: Unit Test
- ./run-fe-ut.sh --run
org.apache.doris.qe.TimeBasedChangeVisibleWaiterTest,org.apache.doris.planner.OlapScanNodeTest
- ./build.sh --fe
- Behavior changed: Yes. Cloud time-based incremental reads wait for
query-start transactions and fail on wait/check timeout or error instead
of succeeding against an incomplete snapshot.
- Does this need documentation: No
---
.../java/org/apache/doris/planner/ScanNode.java | 10 ++-
.../doris/qe/TimeBasedChangeVisibleWaiter.java | 72 +++++++++++++++++++---
.../org/apache/doris/planner/OlapScanNodeTest.java | 35 +++++++++++
.../doris/qe/TimeBasedChangeVisibleWaiterTest.java | 54 ++++++++++++++++
4 files changed, 163 insertions(+), 8 deletions(-)
diff --git a/fe/fe-core/src/main/java/org/apache/doris/planner/ScanNode.java
b/fe/fe-core/src/main/java/org/apache/doris/planner/ScanNode.java
index e5502c0a8a2..efa5d7e5406 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/planner/ScanNode.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/planner/ScanNode.java
@@ -649,6 +649,7 @@ public abstract class ScanNode extends PlanNode implements
SplitGenerator {
List<CloudPartition> partitions = new ArrayList<>();
Set<Long> partitionSet = new HashSet<>();
+ boolean hasIncrementalRead = false;
for (ScanNode node : scanNodes) {
if (!(node instanceof OlapScanNode)) {
continue;
@@ -660,6 +661,9 @@ public abstract class ScanNode extends PlanNode implements
SplitGenerator {
&& ((OlapTableWrapper) table).hasFixedVisibleVersions()) {
continue;
}
+ if (scanNode.getScanParams() != null &&
scanNode.getScanParams().incrementalRead()) {
+ hasIncrementalRead = true;
+ }
for (Long id : scanNode.getSelectedPartitionIds()) {
if (!partitionSet.contains(id)) {
partitionSet.add(id);
@@ -672,7 +676,11 @@ public abstract class ScanNode extends PlanNode implements
SplitGenerator {
if (!partitions.isEmpty()) {
List<Long> versions;
try {
- versions =
CloudPartition.getSnapshotVisibleVersion(partitions);
+ // A time-based change read may have just waited for an old
transaction to finish.
+ // Bypass the FE cache so the scan uses the version made
visible by that transaction.
+ versions = hasIncrementalRead
+ ?
CloudPartition.getSnapshotVisibleVersionFromMs(partitions, false)
+ : CloudPartition.getSnapshotVisibleVersion(partitions);
} catch (RpcException e) {
throw new UserException("get visible version for OlapScanNode
failed", e);
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/qe/TimeBasedChangeVisibleWaiter.java
b/fe/fe-core/src/main/java/org/apache/doris/qe/TimeBasedChangeVisibleWaiter.java
index 6ce5a1b5ab4..13c2d4bf82f 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/qe/TimeBasedChangeVisibleWaiter.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/qe/TimeBasedChangeVisibleWaiter.java
@@ -21,32 +21,40 @@ import org.apache.doris.analysis.TableScanParams;
import org.apache.doris.catalog.Env;
import org.apache.doris.catalog.OlapTable;
import org.apache.doris.catalog.TableIf;
+import org.apache.doris.common.AnalysisException;
+import org.apache.doris.common.Config;
import org.apache.doris.common.Pair;
import org.apache.doris.common.UserException;
import org.apache.doris.nereids.analyzer.UnboundRelation;
import org.apache.doris.nereids.trees.plans.Plan;
import org.apache.doris.nereids.util.RelationUtil;
import org.apache.doris.planner.OlapScanNode;
+import org.apache.doris.transaction.GlobalTransactionMgrIface;
import org.apache.doris.transaction.TransactionState;
import org.apache.doris.transaction.TransactionStatus;
import org.apache.doris.tso.TSOTimestamp;
import com.google.common.annotations.VisibleForTesting;
+import java.util.ArrayList;
+import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
/**
- * Before executing a time-based incremental read, block until every
transaction that committed at
- * or before the requested read timestamp of the target tables becomes
visible. This guarantees the
- * read sees a complete set of changes up to that time point.
+ * Before executing a time-based incremental read, wait for relevant
target-table transactions to
+ * become visible. In cloud mode, drain transactions registered before a
query-start transaction ID
+ * watermark. In non-cloud mode, wait for committed transactions whose commit
TSO is within the
+ * requested read timestamp.
*
* <p>Skipped entirely when the session enables eventual-consistent change
reads, or when no table
* is involved. Waiting is bounded by session variable {@code
change_visible_timeout_ms}; timing out
* raises a {@link UserException}.
*/
public class TimeBasedChangeVisibleWaiter {
+ private static final long CLOUD_TXN_POLL_INTERVAL_MS = 100;
+
private final ConnectContext context;
public static void waitForVisible(ConnectContext context, Plan plan,
Map<List<String>, TableIf> tables)
@@ -88,15 +96,16 @@ public class TimeBasedChangeVisibleWaiter {
return dbToTableEndTSO;
}
- /**
- * For each db, scan its committed-but-not-visible transactions; whenever
a transaction's commit
- * TSO falls within a target table's endTSO, wait for that transaction to
become visible.
- */
+ /** Wait for relevant transactions using the transaction manager
implementation for the cluster mode. */
private void waitForDbToTableEndTSO(Map<Long, Map<Long, Long>>
dbToTableEndTSO) throws UserException {
if (dbToTableEndTSO.isEmpty()) {
return;
}
long deadlineMs = System.currentTimeMillis() +
context.getSessionVariable().getChangeVisibleTimeoutMs();
+ if (Config.isCloudMode()) {
+ waitForCloudTransactions(dbToTableEndTSO, deadlineMs);
+ return;
+ }
for (Map.Entry<Long, Map<Long, Long>> dbEntry :
dbToTableEndTSO.entrySet()) {
long dbId = dbEntry.getKey();
Map<Long, Long> tableEndTSO = dbEntry.getValue();
@@ -110,6 +119,55 @@ public class TimeBasedChangeVisibleWaiter {
}
}
+ private void waitForCloudTransactions(Map<Long, Map<Long, Long>>
dbToTableEndTSO, long deadlineMs)
+ throws UserException {
+ GlobalTransactionMgrIface txnMgr =
Env.getCurrentGlobalTransactionMgr();
+ long txnIdWatermark;
+ try {
+ // MetaService returns the current maximum transaction ID, while
check_txn_conflict
+ // uses an exclusive upper bound.
+ txnIdWatermark = txnMgr.getNextTransactionId() + 1;
+ } catch (UserException e) {
+ throw new UserException("get transaction id watermark failed for
time-based read", e);
+ }
+
+ for (Map.Entry<Long, Map<Long, Long>> dbEntry :
dbToTableEndTSO.entrySet()) {
+ long dbId = dbEntry.getKey();
+ List<Long> tableIds = new ArrayList<>(dbEntry.getValue().keySet());
+ Collections.sort(tableIds);
+ while (!isPreviousTransactionsFinished(txnMgr, txnIdWatermark,
dbId, tableIds)) {
+ long remainingMs = deadlineMs - System.currentTimeMillis();
+ if (remainingMs <= 0) {
+ throw new UserException(String.format(
+ "timeout waiting previous transactions finish for
time-based read, "
+ + "txnIdWatermark=%d dbId=%d tableIds=%s",
+ txnIdWatermark, dbId, tableIds));
+ }
+ try {
+ Thread.sleep(Math.min(CLOUD_TXN_POLL_INTERVAL_MS,
remainingMs));
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new UserException(String.format(
+ "interrupted while waiting previous transactions
finish for time-based read, "
+ + "txnIdWatermark=%d dbId=%d tableIds=%s",
+ txnIdWatermark, dbId, tableIds), e);
+ }
+ }
+ }
+ }
+
+ private boolean isPreviousTransactionsFinished(GlobalTransactionMgrIface
txnMgr, long txnIdWatermark,
+ long dbId, List<Long> tableIds) throws UserException {
+ try {
+ return txnMgr.isPreviousTransactionsFinished(txnIdWatermark, dbId,
tableIds);
+ } catch (AnalysisException e) {
+ throw new UserException(String.format(
+ "check previous transactions failed for time-based read, "
+ + "txnIdWatermark=%d dbId=%d tableIds=%s",
+ txnIdWatermark, dbId, tableIds), e);
+ }
+ }
+
/**
* Return (tableId, endTSO) if the transaction is COMMITTED and its commit
TSO is within the
* requested endTSO of one of its tables; otherwise null (no need to wait).
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/planner/OlapScanNodeTest.java
b/fe/fe-core/src/test/java/org/apache/doris/planner/OlapScanNodeTest.java
index eaba851f69a..443030ef717 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/planner/OlapScanNodeTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/planner/OlapScanNodeTest.java
@@ -25,6 +25,7 @@ import org.apache.doris.analysis.PartitionValue;
import org.apache.doris.analysis.SlotDescriptor;
import org.apache.doris.analysis.SlotId;
import org.apache.doris.analysis.SlotRef;
+import org.apache.doris.analysis.TableScanParams;
import org.apache.doris.analysis.TupleDescriptor;
import org.apache.doris.analysis.TupleId;
import org.apache.doris.catalog.Column;
@@ -41,6 +42,7 @@ import org.apache.doris.catalog.RangePartitionItem;
import org.apache.doris.catalog.Replica.ReplicaState;
import org.apache.doris.catalog.Tablet;
import org.apache.doris.catalog.info.TableNameInfo;
+import org.apache.doris.cloud.catalog.CloudPartition;
import org.apache.doris.common.AnalysisException;
import org.apache.doris.common.Config;
import org.apache.doris.common.util.DebugPointUtil;
@@ -59,6 +61,7 @@ import com.google.common.collect.Range;
import org.apache.commons.collections4.map.CaseInsensitiveMap;
import org.junit.Assert;
import org.junit.Test;
+import org.mockito.MockedStatic;
import org.mockito.Mockito;
import java.util.Collection;
@@ -294,6 +297,38 @@ public class OlapScanNodeTest {
Assert.assertEquals("p_target,p_after",
scanNode.getSelectedPartitionNamesForExplain());
}
+ @Test
+ public void testIncrementalReadGetsVisibleVersionFromMetaService() throws
Exception {
+ long partitionId = 300L;
+ long visibleVersion = 10L;
+ CloudPartition partition = Mockito.mock(CloudPartition.class);
+ Mockito.when(partition.getId()).thenReturn(partitionId);
+ OlapTable table = Mockito.mock(OlapTable.class);
+ Mockito.when(table.getPartition(partitionId)).thenReturn(partition);
+
+ OlapScanNode scanNode = Mockito.mock(OlapScanNode.class);
+ Mockito.when(scanNode.getOlapTable()).thenReturn(table);
+
Mockito.when(scanNode.getSelectedPartitionIds()).thenReturn(Lists.newArrayList(partitionId));
+ Mockito.when(scanNode.getScanParams()).thenReturn(new TableScanParams(
+ TableScanParams.INCREMENTAL_READ, Collections.emptyMap(),
Collections.emptyList()));
+
+ try (MockedStatic<Config> mockedConfig =
Mockito.mockStatic(Config.class);
+ MockedStatic<CloudPartition> mockedPartition =
Mockito.mockStatic(CloudPartition.class)) {
+ mockedConfig.when(Config::isNotCloudMode).thenReturn(false);
+ mockedPartition.when(() ->
CloudPartition.getSnapshotVisibleVersionFromMs(
+ Mockito.anyList(),
Mockito.eq(false))).thenReturn(Lists.newArrayList(visibleVersion));
+
+
ScanNode.setVisibleVersionForOlapScanNodes(Lists.newArrayList(scanNode));
+
+ mockedPartition.verify(() ->
CloudPartition.getSnapshotVisibleVersionFromMs(
+ Mockito.anyList(), Mockito.eq(false)));
+ mockedPartition.verify(() ->
CloudPartition.getSnapshotVisibleVersion(Mockito.anyList()),
+ Mockito.never());
+ }
+
+
Mockito.verify(scanNode).updateScanRangeVersions(Collections.singletonMap(partitionId,
visibleVersion));
+ }
+
@Test
public void testRuntimeFilterBucketMetadataAttachedOnceAcrossWorkers()
throws Exception {
OlapScanNode scanNode = newBucketPruneScanNode(10L);
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/qe/TimeBasedChangeVisibleWaiterTest.java
b/fe/fe-core/src/test/java/org/apache/doris/qe/TimeBasedChangeVisibleWaiterTest.java
index 2109cad72db..7f67ea6c43c 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/qe/TimeBasedChangeVisibleWaiterTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/qe/TimeBasedChangeVisibleWaiterTest.java
@@ -21,6 +21,9 @@ import org.apache.doris.analysis.TableScanParams;
import org.apache.doris.catalog.Database;
import org.apache.doris.catalog.Env;
import org.apache.doris.catalog.OlapTable;
+import org.apache.doris.common.AnalysisException;
+import org.apache.doris.common.Config;
+import org.apache.doris.common.UserException;
import org.apache.doris.nereids.analyzer.UnboundRelation;
import org.apache.doris.nereids.trees.plans.JoinType;
import org.apache.doris.nereids.trees.plans.Plan;
@@ -100,6 +103,57 @@ public class TimeBasedChangeVisibleWaiterTest {
Mockito.verify(txn,
Mockito.times(1)).waitTransactionVisible(Mockito.anyLong());
}
+ @Test
+ public void testCloudWaitForVisibleUsesTransactionIdWatermark() throws
Exception {
+ ConnectContext context = mockContext();
+ OlapTable table = mockOlapTable(DB_ID, TABLE_ID);
+ GlobalTransactionMgrIface txnMgr =
Mockito.mock(GlobalTransactionMgrIface.class);
+ long currentMaxTxnId = 300L;
+ long txnIdWatermark = currentMaxTxnId + 1;
+
Mockito.when(txnMgr.getNextTransactionId()).thenReturn(currentMaxTxnId);
+ Mockito.when(txnMgr.isPreviousTransactionsFinished(
+ txnIdWatermark, DB_ID,
ImmutableList.of(TABLE_ID))).thenReturn(false, true);
+
+ try (MockedStatic<Config> mockedConfig =
Mockito.mockStatic(Config.class);
+ MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class)) {
+ mockedConfig.when(Config::isCloudMode).thenReturn(true);
+
mockedEnv.when(Env::getCurrentGlobalTransactionMgr).thenReturn(txnMgr);
+
+ TimeBasedChangeVisibleWaiter.waitForVisible(context,
newChangeRelation(1, ImmutableMap.of()),
+ ImmutableMap.of(TABLE_QUALIFIER, table));
+ }
+
+ Mockito.verify(txnMgr, Mockito.times(1)).getNextTransactionId();
+ Mockito.verify(txnMgr,
Mockito.times(2)).isPreviousTransactionsFinished(
+ txnIdWatermark, DB_ID, ImmutableList.of(TABLE_ID));
+ Mockito.verify(txnMgr,
Mockito.never()).getCommittedTransactions(Mockito.anyLong());
+ }
+
+ @Test
+ public void testCloudWaitForVisibleFailsWhenConflictCheckFails() throws
Exception {
+ ConnectContext context = mockContext();
+ OlapTable table = mockOlapTable(DB_ID, TABLE_ID);
+ GlobalTransactionMgrIface txnMgr =
Mockito.mock(GlobalTransactionMgrIface.class);
+ long currentMaxTxnId = 300L;
+ long txnIdWatermark = currentMaxTxnId + 1;
+
Mockito.when(txnMgr.getNextTransactionId()).thenReturn(currentMaxTxnId);
+ Mockito.when(txnMgr.isPreviousTransactionsFinished(
+ txnIdWatermark, DB_ID, ImmutableList.of(TABLE_ID)))
+ .thenThrow(new AnalysisException("check transaction conflict
failed"));
+
+ try (MockedStatic<Config> mockedConfig =
Mockito.mockStatic(Config.class);
+ MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class)) {
+ mockedConfig.when(Config::isCloudMode).thenReturn(true);
+
mockedEnv.when(Env::getCurrentGlobalTransactionMgr).thenReturn(txnMgr);
+
+ UserException exception =
Assertions.assertThrows(UserException.class,
+ () -> TimeBasedChangeVisibleWaiter.waitForVisible(
+ context, newChangeRelation(1, ImmutableMap.of()),
+ ImmutableMap.of(TABLE_QUALIFIER, table)));
+ Assertions.assertTrue(exception.getMessage().contains("check
previous transactions failed"));
+ }
+ }
+
private ConnectContext mockContext() {
ConnectContext context = Mockito.mock(ConnectContext.class);
SessionVariable sessionVariable = new SessionVariable();
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]