This is an automated email from the ASF dual-hosted git repository.
Yukang-Lian 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 87a7f9c0be2 [refactor](stream) Localize Cloud Table Stream rewrite
state handling (#66739)
87a7f9c0be2 is described below
commit 87a7f9c0be237fd180822bce4ca9077b9ac1e332
Author: Jamie <[email protected]>
AuthorDate: Mon Aug 24 15:25:14 2026 +0800
[refactor](stream) Localize Cloud Table Stream rewrite state handling
(#66739)
### What problem does this PR solve?
Issue Number: None
Related PR: #65859
Problem Summary: Cloud Table Stream read state is installed only for the
complete statement plan. The previous implementation expressed that
lifecycle by changing the global rewrite pipeline and by teaching the
generic empty-partition rule about Cloud Stream state. This PR restores
the normal rewrite order, keeps `PruneEmptyPartition` generic, and
localizes the readiness check in the Stream wrapper and normalization
rule. A temporary MV plan keeps all candidate partitions and skips
normalization until read state is available; the final statement plan
then normalizes normally. Non-Cloud behavior is unchanged.
---
.../doris/nereids/jobs/executor/Rewriter.java | 72 +++------------------
.../org/apache/doris/nereids/rules/RuleType.java | 1 -
.../doris/nereids/rules/analysis/BindRelation.java | 10 ++-
.../CloudTableStreamReadStateHook.java} | 45 ++++++-------
.../nereids/rules/rewrite/PruneEmptyPartition.java | 7 --
.../insert/StreamConsumptionInfoExtractor.java | 10 ++-
.../InsertIntoTableCommandTableStreamTest.java | 75 +++++++++++++++++-----
7 files changed, 99 insertions(+), 121 deletions(-)
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/executor/Rewriter.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/executor/Rewriter.java
index 58ad5661afa..162289c3355 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/executor/Rewriter.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/executor/Rewriter.java
@@ -17,7 +17,6 @@
package org.apache.doris.nereids.jobs.executor;
-import org.apache.doris.common.Config;
import org.apache.doris.nereids.CascadesContext;
import org.apache.doris.nereids.jobs.JobContext;
import org.apache.doris.nereids.jobs.rewrite.CostBasedRewriteJob;
@@ -157,7 +156,6 @@ import
org.apache.doris.nereids.rules.rewrite.PushProjectThroughUnion;
import org.apache.doris.nereids.rules.rewrite.RecordPlanForMvPreRewrite;
import org.apache.doris.nereids.rules.rewrite.ReduceAggregateChildOutputRows;
import org.apache.doris.nereids.rules.rewrite.ReorderJoin;
-import org.apache.doris.nereids.rules.rewrite.ResolveCloudTableStreamReadState;
import org.apache.doris.nereids.rules.rewrite.RewriteCteChildren;
import org.apache.doris.nereids.rules.rewrite.RewriteSearchToSlots;
import org.apache.doris.nereids.rules.rewrite.RewriteSimpleAggToConstantRule;
@@ -186,7 +184,6 @@ import
org.apache.doris.nereids.trees.plans.logical.LogicalCatalogRelation;
import org.apache.doris.nereids.trees.plans.logical.LogicalFilter;
import org.apache.doris.nereids.trees.plans.logical.LogicalJoin;
import org.apache.doris.nereids.trees.plans.logical.LogicalLimit;
-import org.apache.doris.nereids.trees.plans.logical.LogicalOlapTableStreamScan;
import org.apache.doris.nereids.trees.plans.logical.LogicalSetOperation;
import org.apache.doris.nereids.trees.plans.logical.LogicalTopN;
import org.apache.doris.nereids.trees.plans.logical.LogicalUnion;
@@ -409,19 +406,14 @@ public class Rewriter extends AbstractBatchJobExecutor {
topDown(new EliminateJoinByUnique())
),
topic("Table/Physical optimization",
- // This is a temporary plan used only for
MV matching. Cloud Table Stream
- // read state and lowering must happen
once on the complete statement plan.
topDown(
new PruneOlapScanPartition(),
new PruneEmptyPartition(),
+ new NormalizeOlapTableStreamScan(),
new PruneFileScanPartition(),
new PushDownFilterIntoSchemaScan()
)
),
- topic("Normalize non-Cloud Table Stream for MV
matching",
- cascadesContext -> Config.isNotCloudMode(),
- topDown(new NormalizeOlapTableStreamScan())
- ),
topic("necessary rules before record mv",
topDown(new LimitSortToTopN()),
topDown(new SplitLimit()),
@@ -496,8 +488,7 @@ public class Rewriter extends AbstractBatchJobExecutor {
)
);
- private static final List<RewriteJob>
CTE_CHILDREN_REWRITE_JOBS_BEFORE_TABLE_PHYSICAL_OPTIMIZATION
- = notTraverseChildrenOf(
+ private static final List<RewriteJob>
CTE_CHILDREN_REWRITE_JOBS_BEFORE_SUB_PATH_PUSH_DOWN = notTraverseChildrenOf(
ImmutableSet.of(LogicalCTEAnchor.class),
() -> jobs(
// before `Subquery unnesting` topic, some correlate slots
should have appeared at LogicalApply.left,
@@ -728,14 +719,8 @@ public class Rewriter extends AbstractBatchJobExecutor {
),
custom(RuleType.PULL_UP_PROJECT_EXPR_UNDER_TOPN,
PullUpProjectExprUnderTopN::new)
- )
- )
- );
-
- private static final List<RewriteJob>
CTE_CHILDREN_REWRITE_JOBS_FROM_TABLE_PHYSICAL_OPTIMIZATION
- = notTraverseChildrenOf(
- ImmutableSet.of(LogicalCTEAnchor.class),
- () -> jobs(
+ ),
+ // TODO: these rules should be implementation rules, and
generate alternative physical plans.
topic("Table/Physical optimization",
cascadesContext ->
cascadesContext.rewritePlanContainsTypes(LogicalCatalogRelation.class),
topDown(
@@ -745,9 +730,8 @@ public class Rewriter extends AbstractBatchJobExecutor {
new
LogicalResultSinkToShortCircuitPointQuery(),
new PruneOlapScanPartition(),
new PruneEmptyPartition(),
- // Stream lowering needs the pruned partitions
and its Cloud read state,
- // and must finish before
OperativeColumnDerive treats stream virtual columns
- // as scan slots.
+ // Stream lowering needs the pruned partitions
and must finish before
+ // OperativeColumnDerive treats stream virtual
columns as scan slots.
new NormalizeOlapTableStreamScan(),
new PruneFileScanPartition(),
new PushDownFilterIntoSchemaScan(),
@@ -803,29 +787,6 @@ public class Rewriter extends AbstractBatchJobExecutor {
)
);
- private static final List<RewriteJob>
CLOUD_TABLE_STREAM_READ_STATE_PREPARATION_JOBS
- = ImmutableList.<RewriteJob>builder()
-
.addAll(CTE_CHILDREN_REWRITE_JOBS_BEFORE_TABLE_PHYSICAL_OPTIMIZATION)
- .addAll(notTraverseChildrenOf(
- ImmutableSet.of(LogicalCTEAnchor.class),
- () -> jobs(
- topic("Prune Cloud Table Stream partitions before
resolving read state",
- cascadesContext ->
cascadesContext.rewritePlanContainsTypes(
- LogicalCatalogRelation.class),
- topDown(
- new
LogicalResultSinkToShortCircuitPointQuery(),
- new PruneOlapScanPartition()
- )
- )
- )))
- .build();
-
- private static final List<RewriteJob>
CTE_CHILDREN_REWRITE_JOBS_BEFORE_SUB_PATH_PUSH_DOWN
- = ImmutableList.<RewriteJob>builder()
-
.addAll(CTE_CHILDREN_REWRITE_JOBS_BEFORE_TABLE_PHYSICAL_OPTIMIZATION)
- .addAll(CTE_CHILDREN_REWRITE_JOBS_FROM_TABLE_PHYSICAL_OPTIMIZATION)
- .build();
-
private static final List<RewriteJob>
CTE_CHILDREN_REWRITE_JOBS_AFTER_SUB_PATH_PUSH_DOWN = notTraverseChildrenOf(
ImmutableSet.of(LogicalCTEAnchor.class),
() -> jobs(
@@ -911,13 +872,12 @@ public class Rewriter extends AbstractBatchJobExecutor {
public static Rewriter getWholeTreeRewriterWithCustomJobs(
CascadesContext cascadesContext, List<RewriteJob> jobs) {
List<RewriteJob> wholeTreeRewriteJobs = getWholeTreeRewriteJobs(
- false, false, ImmutableList.of(), jobs, ImmutableList.of(),
true, false);
+ false, false, jobs, ImmutableList.of(), true, false);
return new Rewriter(cascadesContext, wholeTreeRewriteJobs, true);
}
private static List<RewriteJob> getWholeTreeRewriteJobs(boolean
runCboRules) {
return getWholeTreeRewriteJobs(true, true,
- CLOUD_TABLE_STREAM_READ_STATE_PREPARATION_JOBS,
CTE_CHILDREN_REWRITE_JOBS_BEFORE_SUB_PATH_PUSH_DOWN,
CTE_CHILDREN_REWRITE_JOBS_AFTER_SUB_PATH_PUSH_DOWN,
runCboRules, true);
}
@@ -925,7 +885,6 @@ public class Rewriter extends AbstractBatchJobExecutor {
private static List<RewriteJob> getWholeTreeRewriteJobs(
boolean needSubPathPushDown,
boolean needOrExpansion,
- List<RewriteJob> cloudTableStreamReadStatePreparationJobs,
List<RewriteJob> beforePushDownJobs,
List<RewriteJob> afterPushDownJobs,
boolean runCboRules,
@@ -948,22 +907,7 @@ public class Rewriter extends AbstractBatchJobExecutor {
),
topic("record query tmp plan for mv pre rewrite",
custom(RuleType.RECORD_PLAN_FOR_MV_PRE_REWRITE, RecordPlanForMvPreRewrite::new)
- )));
- if (!cloudTableStreamReadStatePreparationJobs.isEmpty()) {
- rewriteJobs.addAll(jobs(
- topic("prepare Cloud Table Stream read state",
- cascadesContext -> Config.isCloudMode()
- &&
cascadesContext.rewritePlanContainsTypes(
-
LogicalOlapTableStreamScan.class),
- custom(RuleType.REWRITE_CTE_CHILDREN,
- () -> new RewriteCteChildren(
-
cloudTableStreamReadStatePreparationJobs, runCboRules)),
-
custom(RuleType.RESOLVE_CLOUD_TABLE_STREAM_READ_STATE,
-
ResolveCloudTableStreamReadState::new),
- custom(RuleType.CLEAR_CONTEXT_STATUS,
ClearContextStatus::new)
- )));
- }
- rewriteJobs.addAll(jobs(
+ ),
topic("rewrite cte sub-tree before sub path push
down",
custom(RuleType.REWRITE_CTE_CHILDREN,
() -> new
RewriteCteChildren(beforePushDownJobs, runCboRules)
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/RuleType.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/RuleType.java
index 0f120cab96b..4a57e891695 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/RuleType.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/RuleType.java
@@ -427,7 +427,6 @@ public enum RuleType {
DISTINCT_AGGREGATE_SPLIT(RuleTypeClass.REWRITE),
PROCESS_SCALAR_AGG_MUST_USE_MULTI_DISTINCT(RuleTypeClass.REWRITE),
// table stream scan rewrite
- RESOLVE_CLOUD_TABLE_STREAM_READ_STATE(RuleTypeClass.REWRITE),
NORMALIZE_OlAP_TABLE_STREAM_SCAN(RuleTypeClass.REWRITE),
// exploration rules
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/BindRelation.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/BindRelation.java
index a468a9eabca..55216a74592 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/BindRelation.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/BindRelation.java
@@ -40,6 +40,7 @@ import
org.apache.doris.catalog.stream.BaseTableStream.StreamScanType;
import org.apache.doris.catalog.stream.OlapTableStream;
import org.apache.doris.catalog.stream.OlapTableStreamWrapper;
import org.apache.doris.catalog.stream.StreamReadMode;
+import org.apache.doris.common.Config;
import org.apache.doris.common.IdGenerator;
import org.apache.doris.common.Pair;
import org.apache.doris.common.util.Util;
@@ -865,7 +866,8 @@ public class BindRelation extends OneAnalysisRuleFactory {
case TEST_EXTERNAL_TABLE:
return new
LogicalTestScan(unboundRelation.getRelationId(), table,
qualifierWithoutTableName);
case STREAM:
- return makeTableStreamScan(table, unboundRelation,
qualifierWithoutTableName);
+ return makeTableStreamScan(table, unboundRelation,
qualifierWithoutTableName,
+ cascadesContext.getStatementContext());
default:
throw new AnalysisException("Unsupported tableType " +
table.getType());
}
@@ -1015,9 +1017,13 @@ public class BindRelation extends OneAnalysisRuleFactory
{
}).collect(ImmutableList.toImmutableList());
}
- private LogicalPlan makeTableStreamScan(TableIf table, UnboundRelation
unboundRelation, List<String> qualifier)
+ private LogicalPlan makeTableStreamScan(TableIf table, UnboundRelation
unboundRelation, List<String> qualifier,
+ StatementContext statementContext)
throws AnalysisException {
if (table instanceof OlapTableStream) {
+ if (Config.isCloudMode()) {
+
statementContext.addPlannerHook(CloudTableStreamReadStateHook.INSTANCE);
+ }
OlapTableStream olapTableStream = (OlapTableStream) table;
LogicalOlapTableStreamScan scan =
makeOlapTableStreamScan(olapTableStream,
unboundRelation, qualifier);
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/ResolveCloudTableStreamReadState.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CloudTableStreamReadStateHook.java
similarity index 81%
rename from
fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/ResolveCloudTableStreamReadState.java
rename to
fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CloudTableStreamReadStateHook.java
index 0a9fa6d5466..dfa36768871 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/ResolveCloudTableStreamReadState.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CloudTableStreamReadStateHook.java
@@ -15,18 +15,17 @@
// specific language governing permissions and limitations
// under the License.
-package org.apache.doris.nereids.rules.rewrite;
+package org.apache.doris.nereids.rules.analysis;
import org.apache.doris.catalog.stream.OlapTableStreamWrapper;
import org.apache.doris.cloud.proto.Cloud;
import org.apache.doris.cloud.rpc.CloudTableStreamReadStateHelper;
-import org.apache.doris.common.Config;
import org.apache.doris.common.UserException;
+import org.apache.doris.nereids.NereidsPlanner;
+import org.apache.doris.nereids.PlannerHook;
import org.apache.doris.nereids.exceptions.AnalysisException;
-import org.apache.doris.nereids.jobs.JobContext;
import org.apache.doris.nereids.trees.plans.Plan;
import org.apache.doris.nereids.trees.plans.logical.LogicalOlapTableStreamScan;
-import org.apache.doris.nereids.trees.plans.visitor.CustomRewriter;
import com.google.common.base.Preconditions;
import com.google.common.collect.ImmutableMap;
@@ -39,23 +38,24 @@ import java.util.List;
import java.util.Map;
import java.util.Set;
-/** Resolve one statement-level read snapshot for every Cloud Table Stream
scan. */
-public class ResolveCloudTableStreamReadState implements CustomRewriter {
+/** Resolves one read snapshot for every Cloud Table Stream relation after
statement analysis. */
+public class CloudTableStreamReadStateHook implements PlannerHook {
+ public static final CloudTableStreamReadStateHook INSTANCE = new
CloudTableStreamReadStateHook();
+
+ private CloudTableStreamReadStateHook() {
+ }
@Override
- public Plan rewriteRoot(Plan plan, JobContext jobContext) {
- if (Config.isNotCloudMode()) {
- return plan;
- }
+ public void afterAnalyze(NereidsPlanner planner) {
+ resolve(planner.getCascadesContext().getRewritePlan());
+ }
+ static void resolve(Plan plan) {
List<LogicalOlapTableStreamScan> scans = new ArrayList<>();
-
plan.collectToList(LogicalOlapTableStreamScan.class::isInstance).forEach(node
-> {
- LogicalOlapTableStreamScan scan = (LogicalOlapTableStreamScan)
node;
- scans.add(scan);
- });
- if (scans.isEmpty()) {
- return plan;
- }
+
plan.collectToList(LogicalOlapTableStreamScan.class::isInstance).forEach(node ->
+ scans.add((LogicalOlapTableStreamScan) node));
+ Preconditions.checkState(!scans.isEmpty(),
+ "Cloud Table Stream read-state hook requires at least one
Stream scan");
boolean readStatesInstalled =
scans.get(0).getTable().hasCloudReadStates();
for (LogicalOlapTableStreamScan scan : scans) {
@@ -69,7 +69,7 @@ public class ResolveCloudTableStreamReadState implements
CustomRewriter {
}
}
if (readStatesInstalled) {
- return plan;
+ return;
}
Map<Cloud.TableStreamIdentityPB, Set<Long>> requestedPartitions = new
LinkedHashMap<>();
@@ -88,7 +88,7 @@ public class ResolveCloudTableStreamReadState implements
CustomRewriter {
if (requestedPartitions.isEmpty()) {
scans.forEach(scan ->
scan.getTable().installCloudReadStates(ImmutableMap.of()));
- return plan;
+ return;
}
Map<Cloud.TableStreamIdentityPB, Map<Long,
Cloud.TableStreamPartitionReadStatePB>> readStates;
@@ -102,10 +102,8 @@ public class ResolveCloudTableStreamReadState implements
CustomRewriter {
OlapTableStreamWrapper wrapper = wrapperEntry.getKey();
Map<Long, Cloud.TableStreamPartitionReadStatePB> bindingStates =
readStates.get(wrapper.getCloudIdentity());
- if (bindingStates == null) {
- wrapper.installCloudReadStates(ImmutableMap.of());
- continue;
- }
+ Preconditions.checkNotNull(bindingStates,
+ "Cloud Table Stream read state is missing for a requested
binding");
ImmutableMap.Builder<Long, Cloud.TableStreamPartitionReadStatePB>
wrapperStates =
ImmutableMap.builderWithExpectedSize(wrapperEntry.getValue().size());
for (long partitionId : wrapperEntry.getValue()) {
@@ -113,6 +111,5 @@ public class ResolveCloudTableStreamReadState implements
CustomRewriter {
}
wrapper.installCloudReadStates(wrapperStates.build());
}
- return plan;
}
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/PruneEmptyPartition.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/PruneEmptyPartition.java
index 5c022615844..aa0fa7fed52 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/PruneEmptyPartition.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/PruneEmptyPartition.java
@@ -19,8 +19,6 @@ package org.apache.doris.nereids.rules.rewrite;
import org.apache.doris.catalog.OlapTable;
import org.apache.doris.catalog.Partition;
-import org.apache.doris.catalog.stream.OlapTableStreamWrapper;
-import org.apache.doris.common.Config;
import org.apache.doris.nereids.rules.Rule;
import org.apache.doris.nereids.rules.RuleType;
import org.apache.doris.nereids.trees.plans.logical.LogicalEmptyRelation;
@@ -50,11 +48,6 @@ public class PruneEmptyPartition extends
OneRewriteRuleFactory {
}
LogicalOlapScan scan = ctx.root;
OlapTable table = scan.getTable();
- if (Config.isCloudMode() && table instanceof OlapTableStreamWrapper
- && !((OlapTableStreamWrapper) table).hasCloudReadStates())
{
- // Cloud Stream emptiness is defined by the statement read
snapshot, not FE-local offsets.
- return null;
- }
List<Long> partitionIdsToPrune = scan.getSelectedPartitionIds();
List<Long> ids =
table.selectNonEmptyPartitionIds(partitionIdsToPrune, scan.getStreamReadMode());
if (ctx.connectContext != null && ctx.connectContext.isTxnModel())
{
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/StreamConsumptionInfoExtractor.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/StreamConsumptionInfoExtractor.java
index ee8388796ae..bb9d7d27610 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/StreamConsumptionInfoExtractor.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/StreamConsumptionInfoExtractor.java
@@ -28,6 +28,7 @@ import org.apache.doris.common.Pair;
import org.apache.doris.nereids.trees.plans.Plan;
import org.apache.doris.nereids.trees.plans.logical.LogicalOlapTableStreamScan;
+import com.google.common.base.Preconditions;
import com.google.common.collect.Maps;
import java.util.ArrayList;
@@ -55,12 +56,9 @@ public class StreamConsumptionInfoExtractor {
LogicalOlapTableStreamScan streamScan =
(LogicalOlapTableStreamScan) scan;
if (!streamScan.isSnapshot()) {
OlapTableStreamWrapper wrapper = streamScan.getTable();
- // The analyzed plan retains scans that the Cloud
rewrite later eliminates, for example
- // `WHERE FALSE`. No partition was read in that case,
so there is no Offset to advance.
- if (Config.isCloudMode() &&
!wrapper.hasCloudReadStates()) {
- return;
- }
- AbstractTableStreamUpdate update =
wrapper.hasCloudReadStates()
+ Preconditions.checkState(Config.isNotCloudMode() ||
wrapper.hasCloudReadStates(),
+ "Cloud Table Stream read state must be
installed during relation analysis");
+ AbstractTableStreamUpdate update = Config.isCloudMode()
? toCloudOlapTableStreamUpdate(wrapper,
streamScan.getSelectedPartitionIds())
: toOlapTableStreamUpdate(wrapper);
if (hasUpdates(update)) {
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertIntoTableCommandTableStreamTest.java
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertIntoTableCommandTableStreamTest.java
index 1ce1e1fdcfc..1c63fe702a2 100755
---
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertIntoTableCommandTableStreamTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertIntoTableCommandTableStreamTest.java
@@ -37,8 +37,9 @@ import org.apache.doris.cloud.rpc.MetaServiceProxy;
import org.apache.doris.common.Config;
import org.apache.doris.common.FeConstants;
import org.apache.doris.common.jmockit.Deencapsulation;
+import org.apache.doris.nereids.NereidsPlanner;
import org.apache.doris.nereids.parser.NereidsParser;
-import org.apache.doris.nereids.rules.rewrite.ResolveCloudTableStreamReadState;
+import org.apache.doris.nereids.rules.analysis.CloudTableStreamReadStateHook;
import org.apache.doris.nereids.trees.plans.Plan;
import org.apache.doris.nereids.trees.plans.logical.LogicalOlapTableStreamScan;
import org.apache.doris.nereids.trees.plans.logical.LogicalPlan;
@@ -460,15 +461,7 @@ public class InsertIntoTableCommandTableStreamTest extends
TestWithFeService {
.setVisibleVersion(8)))
.build();
- String sql = "insert into test_stream.tbl_target "
- + "select * from test_stream.s1 partition (p1)";
- InsertIntoTableCommand command = (InsertIntoTableCommand)
parser.parseSingle(sql);
- connectContext.setStartTime();
- UUID uuid = UUID.randomUUID();
- connectContext.setQueryId(new TUniqueId(uuid.getMostSignificantBits(),
uuid.getLeastSignificantBits()));
- command.initPlan(connectContext, new StmtExecutor(connectContext,
sql), false);
- Plan analyzedPlan = command.getLineagePlan().orElseThrow();
-
+ String sql = "select * from test_stream.s1 partition (p1)";
String previousCloudUniqueId = Config.cloud_unique_id;
String previousMetaServiceEndpoint = Config.meta_service_endpoint;
MetaServiceProxy proxy = Mockito.mock(MetaServiceProxy.class);
@@ -479,7 +472,17 @@ public class InsertIntoTableCommandTableStreamTest extends
TestWithFeService {
mockedProxy.when(MetaServiceProxy::getInstance).thenReturn(proxy);
Mockito.when(proxy.getTableStreamOffset(Mockito.any())).thenReturn(response);
- new
ResolveCloudTableStreamReadState().rewriteRoot(analyzedPlan, null);
+ connectContext.setStartTime();
+ UUID uuid = UUID.randomUUID();
+ connectContext.setQueryId(
+ new TUniqueId(uuid.getMostSignificantBits(),
uuid.getLeastSignificantBits()));
+ PlanChecker checker =
PlanChecker.from(connectContext).analyze(sql);
+
Assertions.assertTrue(checker.getCascadesContext().getStatementContext().getPlannerHooks()
+ .contains(CloudTableStreamReadStateHook.INSTANCE));
+ NereidsPlanner planner = Mockito.mock(NereidsPlanner.class);
+
Mockito.when(planner.getCascadesContext()).thenReturn(checker.getCascadesContext());
+ CloudTableStreamReadStateHook.INSTANCE.afterAnalyze(planner);
+ Plan analyzedPlan = checker.getPlan();
List<TableStreamUpdateInfo> streamUpdateInfos =
StreamConsumptionInfoExtractor.extract(analyzedPlan);
Assertions.assertEquals(1, streamUpdateInfos.size());
@@ -555,13 +558,14 @@ public class InsertIntoTableCommandTableStreamTest
extends TestWithFeService {
connectContext.setQueryId(
new TUniqueId(uuid.getMostSignificantBits(),
uuid.getLeastSignificantBits()));
PlanChecker checker =
PlanChecker.from(connectContext).analyze(sql);
+ Deencapsulation.invoke(CloudTableStreamReadStateHook.class,
"resolve", checker.getPlan());
checker.getCascadesContext().getStatementContext().setForceRecordTmpPlan(true);
checker.rewrite();
List<Plan> tmpPlans =
checker.getCascadesContext().getStatementContext()
.getTmpPlanForMvRewrite();
Assertions.assertFalse(tmpPlans.isEmpty());
- Assertions.assertTrue(tmpPlans.stream().anyMatch(tmpPlan ->
!tmpPlan
+ Assertions.assertTrue(tmpPlans.stream().allMatch(tmpPlan ->
tmpPlan
.collectToList(LogicalOlapTableStreamScan.class::isInstance).isEmpty()));
Assertions.assertTrue(checker.getCascadesContext().getRewritePlan()
.collectToList(LogicalOlapTableStreamScan.class::isInstance).isEmpty());
@@ -590,8 +594,9 @@ public class InsertIntoTableCommandTableStreamTest extends
TestWithFeService {
OlapTable baseTable = (OlapTable)
db.getTableOrMetaException("tbl_stream_base");
long p1 = baseTable.getPartition("p1").getId();
long p2 = baseTable.getPartition("p2").getId();
- String sql = "select k1, k2 from test_stream.s1 where k1 < 100 union
all "
- + "select k1, k2 from test_stream.s1@snapshot() where k1 >=
100 and k1 < 200";
+ String sql = "select k1, k2 from test_stream.s1 partition (p1) where
k1 < 100 union all "
+ + "select k1, k2 from test_stream.s1@snapshot() partition (p2)
"
+ + "where k1 >= 100 and k1 < 200";
String previousCloudUniqueId = Config.cloud_unique_id;
String previousMetaServiceEndpoint = Config.meta_service_endpoint;
@@ -630,6 +635,7 @@ public class InsertIntoTableCommandTableStreamTest extends
TestWithFeService {
new TUniqueId(uuid.getMostSignificantBits(),
uuid.getLeastSignificantBits()));
PlanChecker checker =
PlanChecker.from(connectContext).analyze(sql);
Plan analyzedPlan = checker.getPlan();
+ Deencapsulation.invoke(CloudTableStreamReadStateHook.class,
"resolve", analyzedPlan);
checker.rewrite();
List<TableStreamUpdateInfo> updates =
StreamConsumptionInfoExtractor.extract(analyzedPlan);
@@ -651,7 +657,11 @@ public class InsertIntoTableCommandTableStreamTest extends
TestWithFeService {
}
@Test
- public void testCloudEliminatedStreamScanDoesNotAdvanceOffset() throws
Exception {
+ public void testCloudPredicateEliminationDoesNotShrinkOffsetRange() throws
Exception {
+ Database db = (Database)
Env.getCurrentInternalCatalog().getDbOrMetaException("test_stream");
+ OlapTable baseTable = (OlapTable)
db.getTableOrMetaException("tbl_stream_base");
+ long p1 = baseTable.getPartition("p1").getId();
+ long p2 = baseTable.getPartition("p2").getId();
String sql = "select k1, k2 from test_stream.s1 where false";
String previousCloudUniqueId = Config.cloud_unique_id;
String previousMetaServiceEndpoint = Config.meta_service_endpoint;
@@ -661,6 +671,28 @@ public class InsertIntoTableCommandTableStreamTest extends
TestWithFeService {
Config.meta_service_endpoint = "127.0.0.1:20121";
try (MockedStatic<MetaServiceProxy> mockedProxy =
Mockito.mockStatic(MetaServiceProxy.class)) {
mockedProxy.when(MetaServiceProxy::getInstance).thenReturn(proxy);
+
Mockito.when(proxy.getTableStreamOffset(Mockito.any())).thenAnswer(invocation
-> {
+ Cloud.GetTableStreamOffsetRequest request =
invocation.getArgument(0);
+ Cloud.GetTableStreamOffsetResponse.Builder response =
+ Cloud.GetTableStreamOffsetResponse.newBuilder()
+
.setStatus(Cloud.MetaServiceResponseStatus.newBuilder()
+
.setCode(Cloud.MetaServiceCode.OK));
+ for (Cloud.TableStreamPartitionSetPB binding :
request.getBindingsList()) {
+ Cloud.TableStreamReadBindingResultPB.Builder result =
+
Cloud.TableStreamReadBindingResultPB.newBuilder()
+ .setIdentity(binding.getIdentity());
+ for (long partitionId : binding.getPartitionIdsList())
{
+
result.addPartitionStates(Cloud.TableStreamPartitionReadStatePB.newBuilder()
+ .setPartitionId(partitionId)
+
.setOffsetState(Cloud.TableStreamOffsetStatePB.TABLE_STREAM_OFFSET_CONSUMED)
+ .setOffsetTso(100)
+ .setEndTso(130)
+ .setVisibleVersion(8));
+ }
+ response.addBindings(result);
+ }
+ return response.build();
+ });
connectContext.setStartTime();
UUID uuid = UUID.randomUUID();
@@ -668,10 +700,19 @@ public class InsertIntoTableCommandTableStreamTest
extends TestWithFeService {
new TUniqueId(uuid.getMostSignificantBits(),
uuid.getLeastSignificantBits()));
PlanChecker checker =
PlanChecker.from(connectContext).analyze(sql);
Plan analyzedPlan = checker.getPlan();
+ Deencapsulation.invoke(CloudTableStreamReadStateHook.class,
"resolve", analyzedPlan);
checker.rewrite();
-
Assertions.assertTrue(StreamConsumptionInfoExtractor.extract(analyzedPlan).isEmpty());
- Mockito.verify(proxy,
Mockito.never()).getTableStreamOffset(Mockito.any());
+ List<TableStreamUpdateInfo> updates =
StreamConsumptionInfoExtractor.extract(analyzedPlan);
+ Assertions.assertEquals(1, updates.size());
+ CloudOlapTableStreamUpdate update =
(CloudOlapTableStreamUpdate) updates.get(0).getUpdate();
+ Assertions.assertEquals(Set.of(p1, p2),
update.getPartitionUpdates().keySet());
+
+ ArgumentCaptor<Cloud.GetTableStreamOffsetRequest>
requestCaptor =
+
ArgumentCaptor.forClass(Cloud.GetTableStreamOffsetRequest.class);
+
Mockito.verify(proxy).getTableStreamOffset(requestCaptor.capture());
+ Assertions.assertEquals(Set.of(p1, p2),
+ new
HashSet<>(requestCaptor.getValue().getBindings(0).getPartitionIdsList()));
}
} finally {
Config.cloud_unique_id = previousCloudUniqueId;
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]