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]

Reply via email to