This is an automated email from the ASF dual-hosted git repository.
morrySnow 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 9d576e9e00a [improvement](planner) Reduce planner overhead (#67797)
9d576e9e00a is described below
commit 9d576e9e00a68dadbc76bd3f1d9b566bb6bc71d9
Author: morrySnow <[email protected]>
AuthorDate: Tue Sep 15 16:26:56 2026 +0800
[improvement](planner) Reduce planner overhead (#67797)
### What problem does this PR solve?
Problem Summary: Nereids created ProcessState and maintained
rewrite-path state even when plan-process tracing was disabled, rendered
the final physical plan for every SQL-cache candidate before cache
admission, and repeated cost calculations and CostWeight construction in
the Cascades hot path.
This PR:
- Creates ProcessState only while plan-process tracing is active.
EXPLAIN PLAN PROCESS remains unchanged.
- Defers SQL-cache physical-plan rendering until FE or BE
cache-admission checks succeed.
- Removes redundant pre-regulation node-cost calculation and child-cost
accumulation, while retaining the final property-aware cost
recalculation.
- Keeps CostWeight as the semantic weight snapshot, creates it lazily
once per StatementContext, and reuses it for all Cost construction. The
snapshot is taken after SET_VAR preprocessing and refreshed when a
prepared statement starts another execution.
- Adds immutable Cost addition so accumulated costs do not need to be
reweighted.
The independently measured ProcessState and cost-cleanup budgets are:
| Query | CPU saving | Planning latency saving |
Allocation saving |
| ------------------------- | --------------: | ----------------------: |
----------------: |
| TPCH Q5 | 0.303 ms / 2.7% | 0.552 ms / 4.5% |
200.1 KiB / 4.6% |
| TPCDS Q72 | 0.335 ms / 2.0% | 0.490 ms / 3.0% |
526.3 KiB / 5.4% |
| TPCDS Q64 forced Cascades | 1.712 ms / 5.1% | 2.027 ms / 5.7% |
2618.5 KiB / 9.7% |
These are arithmetic budgets from independently measured constituents,
not a combined-patch ABBA result. Deferred SQL-plan rendering exposes an
additional CPU/allocation cost pool of 0.408 ms/242.8 KiB, 0.588
ms/459.8 KiB, and 1.348 ms/1012.5 KiB respectively. Realized total CPU
savings for that part scale with the non-admission ratio.
---
.../org/apache/doris/nereids/CascadesContext.java | 2 +-
.../org/apache/doris/nereids/NereidsPlanner.java | 22 ++++++--
.../java/org/apache/doris/nereids/PlanContext.java | 13 +++++
.../org/apache/doris/nereids/StatementContext.java | 14 +++++
.../java/org/apache/doris/nereids/cost/Cost.java | 21 +++++---
.../apache/doris/nereids/cost/CostCalculator.java | 13 ++---
.../org/apache/doris/nereids/cost/CostModel.java | 62 ++++++++++------------
.../nereids/jobs/cascades/CostAndEnforcerJob.java | 41 ++++----------
.../jobs/rewrite/BottomUpVisitorRewriteJob.java | 15 +++---
.../jobs/rewrite/TopDownVisitorRewriteJob.java | 12 +++--
.../java/org/apache/doris/nereids/memo/Memo.java | 18 ++++---
.../doris/nereids/minidump/MinidumpUtils.java | 10 ++--
.../properties/ChildrenPropertiesRegulator.java | 8 +--
.../properties/EnforceMissingPropertiesHelper.java | 14 +++--
.../nereids/stats/MemoStatsAndCostRecomputer.java | 7 +--
.../nereids/trees/plans/ComputeResultSet.java | 20 +------
.../plans/physical/PhysicalBlackholeSink.java | 6 +--
.../plans/physical/PhysicalEmptyRelation.java | 18 +------
.../plans/physical/PhysicalOneRowRelation.java | 16 +-----
.../trees/plans/physical/PhysicalResultSink.java | 6 +--
.../trees/plans/physical/PhysicalSqlCache.java | 4 +-
.../apache/doris/nereids/cost/CostModelV1Test.java | 48 +++++++++++++++++
.../doris/nereids/minidump/MinidumpUtTest.java | 18 ++++---
.../doris/nereids/minidump/MinidumpUtTestData.json | 3 +-
.../ChildrenPropertiesRegulatorTest.java | 19 +++----
.../java/org/apache/doris/qe/SqlCacheTest.java | 4 ++
26 files changed, 233 insertions(+), 201 deletions(-)
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/CascadesContext.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/CascadesContext.java
index 4806ebc4c66..1c7ce0fe8be 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/nereids/CascadesContext.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/CascadesContext.java
@@ -271,7 +271,7 @@ public class CascadesContext implements ScheduleContext {
* Init memo with plan
*/
public void toMemo() {
- this.memo = new Memo(getConnectContext(), plan);
+ this.memo = new Memo(getConnectContext(), plan,
statementContext.getCostWeight());
List<Plan> rewrittenPlansByMv =
this.getStatementContext().getRewrittenPlansByMv();
if (!statementContext.getRewrittenPlansByMv().isEmpty()) {
// copy tmp plan for mv rewrite firstly
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/NereidsPlanner.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/NereidsPlanner.java
index 1053699fff0..eae77ce489b 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/nereids/NereidsPlanner.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/NereidsPlanner.java
@@ -88,6 +88,7 @@ import org.apache.doris.qe.ResultSet;
import org.apache.doris.qe.SessionVariable;
import org.apache.doris.qe.TimeBasedChangeVisibleWaiter;
import org.apache.doris.qe.VariableMgr;
+import org.apache.doris.qe.cache.CacheAnalyzer;
import org.apache.doris.statistics.query.QueryStatsRecorder;
import org.apache.doris.statistics.util.StatisticsUtil;
import org.apache.doris.thrift.TQueryCacheParam;
@@ -808,7 +809,7 @@ public class NereidsPlanner extends Planner {
&& cascadesContext.getConnectContext().supportHandleByFe()
&& physicalPlan instanceof ComputeResultSet) {
Optional<ResultSet> resultSet = ((ComputeResultSet)
physicalPlan).computeResultInFe(
- cascadesContext, Optional.empty(),
physicalPlan.getOutput());
+ cascadesContext, physicalPlan.getOutput());
if (resultSet.isPresent()) {
notNeedBackend = true;
}
@@ -1199,10 +1200,10 @@ public class NereidsPlanner extends Planner {
setFormatOptions();
if (physicalPlan instanceof ComputeResultSet) {
- Optional<SqlCacheContext> sqlCacheContext =
statementContext.getSqlCacheContext();
Optional<ResultSet> resultSet = ((ComputeResultSet) physicalPlan)
- .computeResultInFe(cascadesContext, sqlCacheContext,
physicalPlan.getOutput());
+ .computeResultInFe(cascadesContext,
physicalPlan.getOutput());
if (resultSet.isPresent()) {
+ tryAddResultSetToSqlCache(resultSet.get());
return resultSet;
}
}
@@ -1210,6 +1211,21 @@ public class NereidsPlanner extends Planner {
return Optional.empty();
}
+ private void tryAddResultSetToSqlCache(ResultSet resultSet) {
+ if (physicalPlan instanceof PhysicalSqlCache) {
+ return;
+ }
+ Optional<SqlCacheContext> sqlCacheContext =
statementContext.getSqlCacheContext();
+ if (!sqlCacheContext.isPresent()
+ ||
!CacheAnalyzer.canUseSqlCache(statementContext.getConnectContext().getSessionVariable()))
{
+ return;
+ }
+ sqlCacheContext.get().setResultSetInFe(resultSet);
+ Env.getCurrentEnv().getSqlCacheManager().tryAddFeSqlCache(
+ statementContext.getConnectContext(),
+ statementContext.getOriginStatement().originStmt);
+ }
+
private void setFormatOptions() {
ConnectContext ctx = statementContext.getConnectContext();
SessionVariable sessionVariable = ctx.getSessionVariable();
diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/PlanContext.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/PlanContext.java
index 8e9ab3e4407..98730071915 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/nereids/PlanContext.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/PlanContext.java
@@ -17,6 +17,7 @@
package org.apache.doris.nereids;
+import org.apache.doris.nereids.cost.CostWeight;
import org.apache.doris.nereids.memo.GroupExpression;
import org.apache.doris.nereids.properties.DistributionSpecReplicated;
import org.apache.doris.nereids.properties.PhysicalProperties;
@@ -37,14 +38,22 @@ public class PlanContext {
private final ConnectContext connectContext;
private final GroupExpression groupExpression;
private final boolean isBroadcastJoin;
+ private final CostWeight costWeight;
/**
* Constructor for PlanContext.
*/
public PlanContext(ConnectContext connectContext, GroupExpression
groupExpression,
List<PhysicalProperties> childrenProperties) {
+ this(connectContext, groupExpression, childrenProperties, null);
+ }
+
+ /** Constructor for cost calculation with the statement-scoped weight
snapshot. */
+ public PlanContext(ConnectContext connectContext, GroupExpression
groupExpression,
+ List<PhysicalProperties> childrenProperties, CostWeight
costWeight) {
this.connectContext = connectContext;
this.groupExpression = groupExpression;
+ this.costWeight = costWeight;
if (childrenProperties.size() >= 2
&& childrenProperties.get(1).getDistributionSpec() instanceof
DistributionSpecReplicated) {
isBroadcastJoin = true;
@@ -57,6 +66,10 @@ public class PlanContext {
return connectContext.getSessionVariable();
}
+ public CostWeight getCostWeight() {
+ return costWeight;
+ }
+
public boolean isBroadcastJoin() {
return isBroadcastJoin;
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java
index 64670a7398b..d6ea7ebb976 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java
@@ -42,6 +42,7 @@ import org.apache.doris.foundation.format.FormatOptions;
import org.apache.doris.mtmv.BaseTableInfo;
import org.apache.doris.mtmv.ivm.IvmRewriteContext;
import org.apache.doris.nereids.analyzer.UnboundRelation;
+import org.apache.doris.nereids.cost.CostWeight;
import org.apache.doris.nereids.exceptions.AnalysisException;
import org.apache.doris.nereids.hint.Hint;
import org.apache.doris.nereids.hint.UseMvHint;
@@ -128,6 +129,8 @@ public class StatementContext implements Closeable {
}
private ConnectContext connectContext;
+ // Initialized on first cost calculation so per-query SET_VAR hints have
already taken effect.
+ private CostWeight costWeight;
private Optional<IvmRewriteContext> ivmRewriteContext = Optional.empty();
private final Stopwatch stopwatch = Stopwatch.createUnstarted();
@@ -565,6 +568,9 @@ public class StatementContext implements Closeable {
public void setConnectContext(ConnectContext connectContext) {
this.connectContext = connectContext;
+ // Prepared statements reuse their StatementContext across executions.
Each execution must
+ // capture the weights currently effective in the owning
ConnectContext.
+ this.costWeight = null;
}
public void setHasNondeterministic(boolean hasNondeterministic) {
@@ -579,6 +585,14 @@ public class StatementContext implements Closeable {
return connectContext;
}
+ /** Get the cost weights shared by all cost calculations in this
statement. */
+ public CostWeight getCostWeight() {
+ if (costWeight == null) {
+ costWeight = CostWeight.get(connectContext.getSessionVariable());
+ }
+ return costWeight;
+ }
+
public Optional<IvmRewriteContext> getIvmRewriteContext() {
return ivmRewriteContext;
}
diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/cost/Cost.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/cost/Cost.java
index 3677261a265..a7a49a67f6c 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/nereids/cost/Cost.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/cost/Cost.java
@@ -17,8 +17,6 @@
package org.apache.doris.nereids.cost;
-import org.apache.doris.qe.SessionVariable;
-
/**
* CostV1.
*/
@@ -37,7 +35,7 @@ public class Cost {
/**
* Constructor of CostV1.
*/
- public Cost(SessionVariable sessionVariable, double cpuCost, double
memoryCost, double networkCost) {
+ public Cost(CostWeight costWeight, double cpuCost, double memoryCost,
double networkCost) {
// TODO: fix stats
cpuCost = Double.max(0, cpuCost);
memoryCost = Double.max(0, memoryCost);
@@ -46,7 +44,6 @@ public class Cost {
this.memoryCost = memoryCost;
this.networkCost = networkCost;
- CostWeight costWeight = CostWeight.get(sessionVariable);
this.cost = costWeight.cpuWeight * cpuCost + costWeight.memoryWeight *
memoryCost
+ costWeight.networkWeight * networkCost;
}
@@ -82,12 +79,20 @@ public class Cost {
return cost;
}
- public static Cost of(SessionVariable sessionVariable, double cpuCost,
double maxMemory, double networkCost) {
- return new Cost(sessionVariable, cpuCost, maxMemory, networkCost);
+ public static Cost of(CostWeight costWeight, double cpuCost, double
maxMemory, double networkCost) {
+ return new Cost(costWeight, cpuCost, maxMemory, networkCost);
+ }
+
+ public static Cost ofCpu(CostWeight costWeight, double cpuCost) {
+ return new Cost(costWeight, cpuCost, 0, 0);
}
- public static Cost ofCpu(SessionVariable sessionVariable, double cpuCost) {
- return new Cost(sessionVariable, cpuCost, 0, 0);
+ /** Add another cost and compute the weighted value from the summed
components. */
+ public Cost add(Cost other, CostWeight costWeight) {
+ return new Cost(costWeight,
+ cpuCost + other.cpuCost,
+ memoryCost + other.memoryCost,
+ networkCost + other.networkCost);
}
@Override
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/cost/CostCalculator.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/cost/CostCalculator.java
index 8ec69391691..e21ee32ce13 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/nereids/cost/CostCalculator.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/cost/CostCalculator.java
@@ -20,9 +20,7 @@ package org.apache.doris.nereids.cost;
import org.apache.doris.nereids.PlanContext;
import org.apache.doris.nereids.memo.GroupExpression;
import org.apache.doris.nereids.properties.PhysicalProperties;
-import org.apache.doris.nereids.trees.plans.Plan;
import org.apache.doris.qe.ConnectContext;
-import org.apache.doris.qe.SessionVariable;
import java.util.List;
@@ -36,15 +34,10 @@ public class CostCalculator {
* Calculate cost for groupExpression
*/
public static Cost calculateCost(ConnectContext connectContext,
GroupExpression groupExpression,
- List<PhysicalProperties> childrenProperties) {
- PlanContext planContext = new PlanContext(connectContext,
groupExpression, childrenProperties);
+ List<PhysicalProperties> childrenProperties, CostWeight
costWeight) {
+ PlanContext planContext = new PlanContext(
+ connectContext, groupExpression, childrenProperties,
costWeight);
CostModel costModelV1 = new CostModel(connectContext);
return groupExpression.getPlan().accept(costModelV1, planContext);
}
-
- public static Cost addChildCost(ConnectContext connectContext, Plan plan,
Cost planCost, Cost childCost,
- int index) {
- SessionVariable sessionVariable = connectContext.getSessionVariable();
- return CostModel.addChildCost(sessionVariable, planCost, childCost);
- }
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/cost/CostModel.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/cost/CostModel.java
index 2146a2bffec..4c9d2a0b4ba 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/nereids/cost/CostModel.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/cost/CostModel.java
@@ -112,14 +112,6 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
.getHboPlanStatisticsProvider(), "HboPlanStatisticsProvider is
null");
}
- public static Cost addChildCost(SessionVariable sessionVariable, Cost
planCost, Cost childCost) {
- Preconditions.checkArgument(childCost instanceof Cost && planCost
instanceof Cost);
- return new Cost(sessionVariable,
- childCost.getCpuCost() + planCost.getCpuCost(),
- childCost.getMemoryCost() + planCost.getMemoryCost(),
- childCost.getNetworkCost() + planCost.getNetworkCost());
- }
-
@Override
public Cost visit(Plan plan, PlanContext context) {
return Cost.zero();
@@ -142,7 +134,7 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
if
(useMvHint.get().getUseMvTableColumnMap().containsKey(mvQualifier)) {
useMvHint.get().getUseMvTableColumnMap().put(mvQualifier,
true);
useMvHint.get().setStatus(Hint.HintStatus.SUCCESS);
- return Cost.ofCpu(context.getSessionVariable(),
Double.NEGATIVE_INFINITY);
+ return Cost.ofCpu(context.getCostWeight(),
Double.NEGATIVE_INFINITY);
}
}
if
(table.getIndexMetaByIndexId(physicalOlapScan.getSelectedIndexId())
@@ -156,10 +148,10 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
.containsKey(table.getFullQualifiers())) {
useMvHint.get().getUseMvTableColumnMap().put(table.getFullQualifiers(), true);
useMvHint.get().setStatus(Hint.HintStatus.SUCCESS);
- return Cost.ofCpu(context.getSessionVariable(),
Double.NEGATIVE_INFINITY);
+ return Cost.ofCpu(context.getCostWeight(),
Double.NEGATIVE_INFINITY);
}
}
- return Cost.ofCpu(context.getSessionVariable(), rows - aggMvBonus);
+ return Cost.ofCpu(context.getCostWeight(), rows - aggMvBonus);
}
private Set<Column> getColumnForRangePredicate(Set<Expression>
expressions) {
@@ -209,13 +201,13 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
}
}
}
- return Cost.ofCpu(context.getSessionVariable(),
+ return Cost.ofCpu(context.getCostWeight(),
(filter.getConjuncts().size() - prefixIndexMatched + exprCost)
* filterCostFactor);
}
public Cost visitPhysicalSchemaScan(PhysicalSchemaScan physicalSchemaScan,
PlanContext context) {
Statistics statistics = context.getStatisticsWithCheck();
- return Cost.ofCpu(context.getSessionVariable(),
statistics.getRowCount());
+ return Cost.ofCpu(context.getCostWeight(), statistics.getRowCount());
}
@Override
@@ -223,14 +215,14 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
PhysicalStorageLayerAggregate storageLayerAggregate, PlanContext
context) {
Cost costValue = (Cost)
storageLayerAggregate.getRelation().accept(this, context);
// multiply a factor less than 1, so we can select
PhysicalStorageLayerAggregate as far as possible
- return new Cost(context.getSessionVariable(), costValue.getCpuCost() *
0.7, costValue.getMemoryCost(),
+ return new Cost(context.getCostWeight(), costValue.getCpuCost() * 0.7,
costValue.getMemoryCost(),
costValue.getNetworkCost());
}
@Override
public Cost visitPhysicalFileScan(PhysicalFileScan physicalFileScan,
PlanContext context) {
Statistics statistics = context.getStatisticsWithCheck();
- return Cost.ofCpu(context.getSessionVariable(),
statistics.getRowCount() * EXTERNAL_TABLE_SCAN_FACTOR);
+ return Cost.ofCpu(context.getCostWeight(), statistics.getRowCount() *
EXTERNAL_TABLE_SCAN_FACTOR);
}
@Override
@@ -245,13 +237,13 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
return Cost.zero();
}
double exprCost = expressionTreeCost(physicalProject.getProjects());
- return Cost.ofCpu(context.getSessionVariable(), exprCost + 1);
+ return Cost.ofCpu(context.getCostWeight(), exprCost + 1);
}
@Override
public Cost visitPhysicalOdbcScan(PhysicalOdbcScan physicalOdbcScan,
PlanContext context) {
Statistics statistics = context.getStatisticsWithCheck();
- return Cost.ofCpu(context.getSessionVariable(),
statistics.getRowCount() * EXTERNAL_TABLE_SCAN_FACTOR);
+ return Cost.ofCpu(context.getCostWeight(), statistics.getRowCount() *
EXTERNAL_TABLE_SCAN_FACTOR);
}
@Override
@@ -267,7 +259,7 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
// Now we do more like two-phase sort, so penalise one-phase sort
rowCount *= 100;
}
- return Cost.of(context.getSessionVariable(), childRowCount, rowCount,
childRowCount);
+ return Cost.of(context.getCostWeight(), childRowCount, rowCount,
childRowCount);
}
@Override
@@ -282,14 +274,14 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
// Now we do more like two-phase sort, so penalise one-phase sort
rowCount = rowCount * 100 + 100;
}
- return Cost.of(context.getSessionVariable(), childRowCount, rowCount,
childRowCount);
+ return Cost.of(context.getCostWeight(), childRowCount, rowCount,
childRowCount);
}
@Override
public Cost visitPhysicalPartitionTopN(PhysicalPartitionTopN<? extends
Plan> partitionTopN, PlanContext context) {
Statistics statistics = context.getStatisticsWithCheck();
Statistics childStatistics = context.getChildStatistics(0);
- return Cost.of(context.getSessionVariable(),
+ return Cost.of(context.getCostWeight(),
childStatistics.getRowCount(),
statistics.getRowCount(),
childStatistics.getRowCount());
@@ -306,7 +298,7 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
double dataSizeFactor =
childStatistics.dataSizeFactor(distribute.child().getOutput());
// shuffle
if (spec instanceof DistributionSpecHash) {
- return Cost.of(context.getSessionVariable(),
+ return Cost.of(context.getCostWeight(),
intputRowCount / beNumForDist,
0,
intputRowCount * dataSizeFactor / beNumForDist
@@ -318,7 +310,7 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
// estimate broadcast cost by an experience formula: beNumber^0.5
* rowCount
// - sender number and receiver number is not available at RBO
stage now, so we use beNumber
// - senders and receivers work in parallel, that why we use
square of beNumber
- return Cost.of(context.getSessionVariable(),
+ return Cost.of(context.getCostWeight(),
0,
0,
intputRowCount * dataSizeFactor);
@@ -327,7 +319,7 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
// gather
if (spec instanceof DistributionSpecGather) {
- return Cost.of(context.getSessionVariable(),
+ return Cost.of(context.getCostWeight(),
0,
0,
intputRowCount * dataSizeFactor / beNumForDist);
@@ -335,7 +327,7 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
// any
// cost of random shuffle is lower than hash shuffle.
- return Cost.of(context.getSessionVariable(),
+ return Cost.of(context.getCostWeight(),
0,
0,
intputRowCount * dataSizeFactor
@@ -359,7 +351,7 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
Statistics inputStatistics = context.getChildStatistics(0);
double exprCost = expressionTreeCost(aggregate.getExpressions());
if (aggregate.getAggPhase().isLocal()) {
- return Cost.of(context.getSessionVariable(),
+ return Cost.of(context.getCostWeight(),
exprCost / 100 + inputStatistics.getRowCount() / beNumber,
inputStatistics.getRowCount() / beNumber, 0);
} else {
@@ -376,7 +368,7 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
aggregate.getGroupByExpressions().size())) {
rowCost *= BUCKETED_AGG_COST_DISCOUNT;
}
- return Cost.of(context.getSessionVariable(),
+ return Cost.of(context.getCostWeight(),
exprCost / 100 + rowCost, rowCost, 0);
}
}
@@ -439,7 +431,7 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
in pattern2, join1 and join2 takes more time, but Agg1 and agg2 can be
processed in parallel.
*/
if (physicalHashJoin.getJoinType().isCrossJoin()) {
- return Cost.of(context.getSessionVariable(), leftRowCount +
rightRowCount + outputRowCount,
+ return Cost.of(context.getCostWeight(), leftRowCount +
rightRowCount + outputRowCount,
0,
leftRowCount + rightRowCount
);
@@ -500,14 +492,14 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
}
}
- return Cost.of(context.getSessionVariable(),
+ return Cost.of(context.getCostWeight(),
leftRowCount * probeShortcutFactor + rightRowCount *
probeShortcutFactor * buildSideFactor
+ outputRowCount * probeSideFactor,
rightRowCount,
0
);
}
- return Cost.of(context.getSessionVariable(),
+ return Cost.of(context.getCostWeight(),
leftRowCount * probeShortcutFactor + rightRowCount *
probeShortcutFactor + outputRowCount,
rightRowCount, 0
);
@@ -568,7 +560,7 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
nljPenalty = Math.min(leftStatistics.getRowCount(),
rightStatistics.getRowCount());
}
nljPenalty = Math.max(nljPenalty, 1.0);
- return Cost.of(context.getSessionVariable(),
+ return Cost.of(context.getCostWeight(),
leftStatistics.getRowCount() * rightStatistics.getRowCount() *
nljPenalty,
rightStatistics.getRowCount() * nljPenalty,
0);
@@ -577,7 +569,7 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
@Override
public Cost visitPhysicalAssertNumRows(PhysicalAssertNumRows<? extends
Plan> assertNumRows,
PlanContext context) {
- return Cost.of(context.getSessionVariable(),
+ return Cost.of(context.getCostWeight(),
assertNumRows.getAssertNumRowsElement().getDesiredNumOfRows(),
assertNumRows.getAssertNumRowsElement().getDesiredNumOfRows(),
0
@@ -592,13 +584,13 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
}
double memoryCost = context.getChildStatistics(0).computeSize(
physicalIntersect.child(0).getOutput());
- return Cost.of(context.getSessionVariable(), cpuCost, memoryCost, 0);
+ return Cost.of(context.getCostWeight(), cpuCost, memoryCost, 0);
}
@Override
public Cost visitPhysicalGenerate(PhysicalGenerate<? extends Plan>
generate, PlanContext context) {
Statistics statistics = context.getStatisticsWithCheck();
- return Cost.of(context.getSessionVariable(),
+ return Cost.of(context.getCostWeight(),
statistics.getRowCount(),
statistics.getRowCount(),
0
@@ -635,7 +627,7 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
}
}
}
- return Cost.of(context.getSessionVariable(), rows, 0, rows * tupleSize
* networkFactor);
+ return Cost.of(context.getCostWeight(), rows, 0, rows * tupleSize *
networkFactor);
}
@Override
@@ -665,6 +657,6 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
}
}
}
- return Cost.of(context.getSessionVariable(), rows, 0, rows * tupleSize
* networkFactor);
+ return Cost.of(context.getCostWeight(), rows, 0, rows * tupleSize *
networkFactor);
}
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/cascades/CostAndEnforcerJob.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/cascades/CostAndEnforcerJob.java
index 31206e85e38..c4b0936c34f 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/cascades/CostAndEnforcerJob.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/cascades/CostAndEnforcerJob.java
@@ -20,6 +20,7 @@ package org.apache.doris.nereids.jobs.cascades;
import org.apache.doris.common.Pair;
import org.apache.doris.nereids.cost.Cost;
import org.apache.doris.nereids.cost.CostCalculator;
+import org.apache.doris.nereids.cost.CostWeight;
import org.apache.doris.nereids.jobs.Job;
import org.apache.doris.nereids.jobs.JobContext;
import org.apache.doris.nereids.jobs.JobType;
@@ -54,8 +55,6 @@ public class CostAndEnforcerJob extends Job implements
Cloneable {
// cost of current plan tree
private Cost curTotalCost;
- // cost of current plan node
- private Cost curNodeCost;
// List of request property to children
// Example: Physical Hash Join
@@ -121,8 +120,6 @@ public class CostAndEnforcerJob extends Job implements
Cloneable {
countJobExecutionTimesOfGroupExpressions(groupExpression);
// Do init logic of root plan/groupExpr of `subplan`, only run once
per task.
if (curChildIndex == -1) {
- curNodeCost = Cost.zero();
- curTotalCost = Cost.zero();
curChildIndex = 0;
// List<request property to children>
// [ child item: [leftProperties, rightProperties]]
@@ -142,14 +139,6 @@ public class CostAndEnforcerJob extends Job implements
Cloneable {
=
requestChildrenPropertiesList.get(requestPropertiesIndex);
List<PhysicalProperties> outputChildrenProperties
= outputChildrenPropertiesList.get(requestPropertiesIndex);
- // Calculate cost
- if (curChildIndex == 0 && prevChildIndex == -1) {
- curNodeCost =
CostCalculator.calculateCost(getConnectContext(), groupExpression,
- requestChildrenProperties);
- groupExpression.setCost(curNodeCost);
- curTotalCost = curNodeCost;
- }
-
// Handle all child plan node.
for (; curChildIndex < groupExpression.arity(); curChildIndex++) {
PhysicalProperties requestChildProperty =
requestChildrenProperties.get(curChildIndex);
@@ -188,12 +177,6 @@ public class CostAndEnforcerJob extends Job implements
Cloneable {
// plan's requestChildProperty).getOutputProperties(current
plan's requestChildProperty) == child
// plan's outputProperties`, the outputProperties must satisfy
the origin requestChildProperty
outputChildrenProperties.set(curChildIndex, outputProperties);
- curTotalCost = CostCalculator.addChildCost(
- getConnectContext(),
- groupExpression.getPlan(),
- curNodeCost,
-
lowestCostExpr.getCostValueByProperties(requestChildProperty),
- curChildIndex);
// Not performing lower bound group pruning here is to avoid
redundant optimization of children.
// For example:
@@ -239,6 +222,7 @@ public class CostAndEnforcerJob extends Job implements
Cloneable {
}
boolean hasSuccess = false;
+ CostWeight costWeight =
context.getCascadesContext().getStatementContext().getCostWeight();
for (List<PhysicalProperties> outputChildrenProperties :
childrenOutputSpace) {
// Not need to do pruning here because it has been done when we
get the
// best expr from the child group
@@ -259,17 +243,14 @@ public class CostAndEnforcerJob extends Job implements
Cloneable {
}
// recompute cost after adjusting property
- curNodeCost = CostCalculator.calculateCost(getConnectContext(),
groupExpression, requestChildrenProperties);
- groupExpression.setCost(curNodeCost);
- curTotalCost = curNodeCost;
+ Cost nodeCost = CostCalculator.calculateCost(
+ getConnectContext(), groupExpression,
requestChildrenProperties, costWeight);
+ groupExpression.setCost(nodeCost);
+ curTotalCost = nodeCost;
for (int i = 0; i < outputChildrenProperties.size(); i++) {
PhysicalProperties childProperties =
outputChildrenProperties.get(i);
- curTotalCost = CostCalculator.addChildCost(
- getConnectContext(),
- groupExpression.getPlan(),
- curTotalCost,
-
groupExpression.child(i).getLowestCostPlan(childProperties).get().first,
- i);
+ curTotalCost = curTotalCost.add(
+
groupExpression.child(i).getLowestCostPlan(childProperties).get().first,
costWeight);
}
// record map { outputProperty -> outputProperty }, { ANY ->
outputProperty },
@@ -313,8 +294,9 @@ public class CostAndEnforcerJob extends Job implements
Cloneable {
}
}
- EnforceMissingPropertiesHelper enforceMissingPropertiesHelper
- = new EnforceMissingPropertiesHelper(getConnectContext(),
groupExpression, curTotalCost);
+ CostWeight costWeight =
context.getCascadesContext().getStatementContext().getCostWeight();
+ EnforceMissingPropertiesHelper enforceMissingPropertiesHelper = new
EnforceMissingPropertiesHelper(
+ getConnectContext(), groupExpression, curTotalCost,
costWeight);
PhysicalProperties addEnforcedProperty = enforceMissingPropertiesHelper
.enforceProperty(outputProperty, requiredProperties);
curTotalCost = enforceMissingPropertiesHelper.getCurTotalCost();
@@ -352,7 +334,6 @@ public class CostAndEnforcerJob extends Job implements
Cloneable {
prevChildIndex = -1;
curChildIndex = 0;
curTotalCost = Cost.zero();
- curNodeCost = Cost.zero();
}
/**
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/rewrite/BottomUpVisitorRewriteJob.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/rewrite/BottomUpVisitorRewriteJob.java
index 7bec8269ccf..2f8f257a9ed 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/rewrite/BottomUpVisitorRewriteJob.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/rewrite/BottomUpVisitorRewriteJob.java
@@ -57,9 +57,9 @@ public class BottomUpVisitorRewriteJob implements RewriteJob {
return;
}
- Plan root = rewrite(
- null, -1, originPlan, jobContext, rules,
batchId.incrementAndGet(), false, new ProcessState(originPlan)
- );
+ ProcessState processState =
jobContext.getCascadesContext().showPlanProcess()
+ ? new ProcessState(originPlan) : null;
+ Plan root = rewrite(null, -1, originPlan, jobContext, rules,
batchId.incrementAndGet(), false, processState);
jobContext.getCascadesContext().setRewritePlan(root);
}
@@ -75,9 +75,6 @@ public class BottomUpVisitorRewriteJob implements RewriteJob {
if (state == RewriteState.REWRITTEN) {
return plan;
}
- CascadesContext cascadesContext = jobContext.getCascadesContext();
- boolean showPlanProcess = cascadesContext.showPlanProcess();
-
Plan currentPlan = plan;
while (true) {
if (fastReturn &&
rules.getCurrentAndChildrenRules(currentPlan).isEmpty()) {
@@ -100,14 +97,14 @@ public class BottomUpVisitorRewriteJob implements
RewriteJob {
}
if (changed) {
currentPlan = currentPlan.withChildren(newChildren.build());
- if (showPlanProcess) {
+ if (processState != null) {
parent = processState.updateChild(parent, childIndex,
currentPlan);
}
}
Plan rewrittenPlan = doRewrite(parent, childIndex, currentPlan,
jobContext, rules, processState);
if (!rewrittenPlan.deepEquals(currentPlan)) {
currentPlan = rewrittenPlan;
- if (showPlanProcess) {
+ if (processState != null) {
parent = processState.updateChild(parent, childIndex,
currentPlan);
}
} else {
@@ -132,7 +129,7 @@ public class BottomUpVisitorRewriteJob implements
RewriteJob {
Plan result = transform.get(0);
currentRule.acceptPlan(result);
- if (cascadesContext.showPlanProcess()) {
+ if (processState != null) {
String beforeShape =
processState.getNewestPlan().treeString(true, plan);
String afterShape =
processState.updateChildAndGetNewest(originParent, childIndex, result)
.treeString(true, result);
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/rewrite/TopDownVisitorRewriteJob.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/rewrite/TopDownVisitorRewriteJob.java
index 8ae29a77f47..bb07c113d23 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/rewrite/TopDownVisitorRewriteJob.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/jobs/rewrite/TopDownVisitorRewriteJob.java
@@ -53,9 +53,9 @@ public class TopDownVisitorRewriteJob implements RewriteJob {
return;
}
- Plan root = rewrite(
- null, -1, originPlan, jobContext, rules, false, new
ProcessState(originPlan)
- );
+ ProcessState processState =
jobContext.getCascadesContext().showPlanProcess()
+ ? new ProcessState(originPlan) : null;
+ Plan root = rewrite(null, -1, originPlan, jobContext, rules, false,
processState);
jobContext.getCascadesContext().setRewritePlan(root);
}
@@ -88,7 +88,9 @@ public class TopDownVisitorRewriteJob implements RewriteJob {
if (changed) {
plan = plan.withChildren(newChildren.build());
- processState.updateChild(parent, childIndex, plan);
+ if (processState != null) {
+ processState.updateChild(parent, childIndex, plan);
+ }
}
return plan;
@@ -113,7 +115,7 @@ public class TopDownVisitorRewriteJob implements RewriteJob
{
if (!transform.isEmpty() &&
!transform.get(0).deepEquals(originPlan)) {
Plan newPlan = transform.get(0);
currentRule.acceptPlan(originPlan);
- if (cascadesContext.showPlanProcess()) {
+ if (processState != null) {
String beforeShape =
processState.getNewestPlan().treeString(true, originPlan);
String afterShape =
processState.updateChildAndGetNewest(originParent, childIndex, newPlan)
.treeString(true, newPlan);
diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/memo/Memo.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/memo/Memo.java
index 0d6862a3f4a..7650b597b11 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/nereids/memo/Memo.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/memo/Memo.java
@@ -22,6 +22,7 @@ import org.apache.doris.common.IdGenerator;
import org.apache.doris.common.Pair;
import org.apache.doris.nereids.cost.Cost;
import org.apache.doris.nereids.cost.CostCalculator;
+import org.apache.doris.nereids.cost.CostWeight;
import org.apache.doris.nereids.exceptions.AnalysisException;
import org.apache.doris.nereids.metrics.EventChannel;
import org.apache.doris.nereids.metrics.EventProducer;
@@ -78,6 +79,7 @@ public class Memo {
EventChannel.getDefaultChannel().addConsumers(new
LogConsumer(GroupMergeEvent.class, EventChannel.LOG)));
private static long stateId = 0;
private final ConnectContext connectContext;
+ private final CostWeight costWeight;
// The key is the query tableId, the value is the refresh version when
last refresh, this is needed
// because struct info refresh base on target tableId.
private final Map<Integer, AtomicInteger> refreshVersion = new HashMap<>();
@@ -95,11 +97,18 @@ public class Memo {
public Memo() {
this.root = null;
this.connectContext = null;
+ this.costWeight = null;
}
public Memo(ConnectContext connectContext, Plan plan) {
+ this(connectContext, plan,
+ connectContext == null ? null :
connectContext.getStatementContext().getCostWeight());
+ }
+
+ public Memo(ConnectContext connectContext, Plan plan, CostWeight
costWeight) {
this.root = init(plan);
this.connectContext = connectContext;
+ this.costWeight = costWeight;
}
public static long getStateId() {
@@ -1035,15 +1044,12 @@ public class Memo {
List<Pair<Long, List<Integer>>> childrenId = new ArrayList<>();
permute(children, 0, childrenId, new ArrayList<>());
- Cost cost = CostCalculator.calculateCost(connectContext,
groupExpression, inputProperties);
+ Cost cost = CostCalculator.calculateCost(
+ connectContext, groupExpression, inputProperties,
costWeight);
for (Pair<Long, List<Integer>> c : childrenId) {
Cost totalCost = cost;
for (int i = 0; i < children.size(); i++) {
- totalCost = CostCalculator.addChildCost(connectContext,
- groupExpression.getPlan(),
- totalCost,
- children.get(i).get(c.second.get(i)).second,
- i);
+ totalCost =
totalCost.add(children.get(i).get(c.second.get(i)).second, costWeight);
}
if (res.isEmpty()) {
Preconditions.checkArgument(
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/minidump/MinidumpUtils.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/minidump/MinidumpUtils.java
index c5f2e44c14d..adb509be75c 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/minidump/MinidumpUtils.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/minidump/MinidumpUtils.java
@@ -223,10 +223,12 @@ public class MinidumpUtils {
*/
public static void setConnectContext(Minidump minidump) {
ConnectContext connectContext = new ConnectContext();
+ Env env = Env.getCurrentEnv();
+ connectContext.setEnv(env);
connectContext.getTotalColumnStatisticMap().putAll(minidump.getTotalColumnStatisticMap());
connectContext.getTotalHistogramMap().putAll(minidump.getTotalHistogramMap());
connectContext.setThreadLocalInfo();
-
Env.getCurrentEnv().setColocateTableIndex(minidump.getColocateTableIndex());
+ env.setColocateTableIndex(minidump.getColocateTableIndex());
connectContext.setSessionVariable(minidump.getSessionVariable());
connectContext.setDatabase(minidump.getDbName());
connectContext.getSessionVariable().setPlanNereidsDump(true);
@@ -244,8 +246,10 @@ public class MinidumpUtils {
if (parsed instanceof ExplainCommand) {
parsed = ((ExplainCommand) parsed).getLogicalPlan();
}
- NereidsPlanner nereidsPlanner = new NereidsPlanner(
- new StatementContext(ConnectContext.get(), new
OriginStatement(sql, 0)));
+ ConnectContext connectContext = ConnectContext.get();
+ StatementContext statementContext = new
StatementContext(connectContext, new OriginStatement(sql, 0));
+ connectContext.setStatementContext(statementContext);
+ NereidsPlanner nereidsPlanner = new NereidsPlanner(statementContext);
nereidsPlanner.plan(LogicalPlanAdapter.of(parsed));
return ((AbstractPlan) nereidsPlanner.getOptimizedPlan()).toJson();
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/properties/ChildrenPropertiesRegulator.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/properties/ChildrenPropertiesRegulator.java
index a5bd2a3c1bf..9741d1de68a 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/properties/ChildrenPropertiesRegulator.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/properties/ChildrenPropertiesRegulator.java
@@ -20,6 +20,7 @@ package org.apache.doris.nereids.properties;
import org.apache.doris.common.Pair;
import org.apache.doris.nereids.cost.Cost;
import org.apache.doris.nereids.cost.CostCalculator;
+import org.apache.doris.nereids.cost.CostWeight;
import org.apache.doris.nereids.jobs.JobContext;
import org.apache.doris.nereids.memo.Group;
import org.apache.doris.nereids.memo.GroupExpression;
@@ -992,10 +993,11 @@ public class ChildrenPropertiesRegulator extends
PlanVisitor<List<List<PhysicalP
GroupExpression enforcer = target.addEnforcer(child.getOwnerGroup());
child.getOwnerGroup().addEnforcer(enforcer);
ConnectContext connectContext =
jobContext.getCascadesContext().getConnectContext();
- Cost enforceCost = CostCalculator.calculateCost(connectContext,
enforcer, Lists.newArrayList(childOutput));
+ CostWeight costWeight =
jobContext.getCascadesContext().getStatementContext().getCostWeight();
+ Cost enforceCost = CostCalculator.calculateCost(
+ connectContext, enforcer, Lists.newArrayList(childOutput),
costWeight);
enforcer.setCost(enforceCost);
- Cost totalCost = CostCalculator.addChildCost(
- connectContext, enforcer.getPlan(), enforceCost, currentCost,
0);
+ Cost totalCost = enforceCost.add(currentCost, costWeight);
if (enforcer.updateLowestCostTable(newOutputProperty,
Lists.newArrayList(childOutput), totalCost)) {
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/properties/EnforceMissingPropertiesHelper.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/properties/EnforceMissingPropertiesHelper.java
index 81ce94e5f2f..4b34d5c19fe 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/properties/EnforceMissingPropertiesHelper.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/properties/EnforceMissingPropertiesHelper.java
@@ -19,6 +19,7 @@ package org.apache.doris.nereids.properties;
import org.apache.doris.nereids.cost.Cost;
import org.apache.doris.nereids.cost.CostCalculator;
+import org.apache.doris.nereids.cost.CostWeight;
import org.apache.doris.nereids.memo.GroupExpression;
import org.apache.doris.nereids.metrics.EventChannel;
import org.apache.doris.nereids.metrics.EventProducer;
@@ -40,13 +41,15 @@ public class EnforceMissingPropertiesHelper {
EventChannel.getDefaultChannel().addConsumers(new
LogConsumer(EnforcerEvent.class, EventChannel.LOG)));
private final ConnectContext connectContext;
private final GroupExpression groupExpression;
+ private final CostWeight costWeight;
private Cost curTotalCost;
public EnforceMissingPropertiesHelper(ConnectContext connectContext,
GroupExpression groupExpression,
- Cost curTotalCost) {
+ Cost curTotalCost, CostWeight costWeight) {
this.connectContext = connectContext;
this.groupExpression = groupExpression;
this.curTotalCost = curTotalCost;
+ this.costWeight = costWeight;
}
public Cost getCurTotalCost() {
@@ -168,15 +171,10 @@ public class EnforceMissingPropertiesHelper {
if (enforcerCost == null) {
enforcer.setEstOutputRowCount(enforcer.getOwnerGroup().getStatistics().getRowCount());
enforcerCost = CostCalculator.calculateCost(connectContext,
enforcer,
- Lists.newArrayList(oldOutputProperty));
+ Lists.newArrayList(oldOutputProperty), costWeight);
enforcer.setCost(enforcerCost);
}
- curTotalCost = CostCalculator.addChildCost(
- connectContext,
- enforcer.getPlan(),
- enforcerCost,
- curTotalCost,
- 0);
+ curTotalCost = enforcerCost.add(curTotalCost, costWeight);
if (enforcer.updateLowestCostTable(newOutputProperty,
Lists.newArrayList(oldOutputProperty), curTotalCost)) {
enforcer.putOutputPropertiesMap(newOutputProperty,
newOutputProperty);
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/stats/MemoStatsAndCostRecomputer.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/stats/MemoStatsAndCostRecomputer.java
index 08ce499974f..8b668d2b872 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/stats/MemoStatsAndCostRecomputer.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/stats/MemoStatsAndCostRecomputer.java
@@ -21,6 +21,7 @@ import org.apache.doris.common.Pair;
import org.apache.doris.nereids.CascadesContext;
import org.apache.doris.nereids.cost.Cost;
import org.apache.doris.nereids.cost.CostCalculator;
+import org.apache.doris.nereids.cost.CostWeight;
import org.apache.doris.nereids.memo.Group;
import org.apache.doris.nereids.memo.GroupExpression;
import org.apache.doris.nereids.properties.PhysicalProperties;
@@ -314,6 +315,7 @@ public final class MemoStatsAndCostRecomputer {
groupExpression.clearCostState();
Cost bestNodeCost = null;
+ CostWeight costWeight =
cascadesContext.getStatementContext().getCostWeight();
for (Map.Entry<PhysicalProperties, Pair<Cost,
List<PhysicalProperties>>> entry
: originalLowestCostTable.entrySet()) {
PhysicalProperties outputProperties = entry.getKey();
@@ -322,7 +324,7 @@ public final class MemoStatsAndCostRecomputer {
continue;
}
Cost nodeCost =
CostCalculator.calculateCost(cascadesContext.getConnectContext(),
- groupExpression, childInputProperties);
+ groupExpression, childInputProperties, costWeight);
Cost totalCost = nodeCost;
for (int i = 0; i < childInputProperties.size(); i++) {
Optional<Pair<Cost, GroupExpression>> childBestPlan =
groupExpression.child(i)
@@ -331,8 +333,7 @@ public final class MemoStatsAndCostRecomputer {
totalCost = null;
break;
}
- totalCost =
CostCalculator.addChildCost(cascadesContext.getConnectContext(),
- groupExpression.getPlan(), totalCost,
childBestPlan.get().first, i);
+ totalCost = totalCost.add(childBestPlan.get().first,
costWeight);
}
if (totalCost == null) {
continue;
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/ComputeResultSet.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/ComputeResultSet.java
index f86e143ca7b..c46b01d00a4 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/ComputeResultSet.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/ComputeResultSet.java
@@ -18,7 +18,6 @@
package org.apache.doris.nereids.trees.plans;
import org.apache.doris.nereids.CascadesContext;
-import org.apache.doris.nereids.SqlCacheContext;
import org.apache.doris.nereids.trees.expressions.Slot;
import org.apache.doris.qe.ResultSet;
@@ -36,23 +35,8 @@ import java.util.Optional;
* the PhysicalEmptyRelation implement this interface.
* </li>
* </p>
- * <p>
- * If you want to cache the result set in fe, you can implement this
interface and write this code:
- * </p>
- * <pre>
- * StatementContext statementContext = cascadesContext.getStatementContext();
- * boolean enableSqlCache
- * =
CacheAnalyzer.canUseSqlCache(statementContext.getConnectContext().getSessionVariable());
- * if (sqlCacheContext.isPresent() && enableSqlCache) {
- * sqlCacheContext.get().setResultSetInFe(resultSet);
- * Env.getCurrentEnv().getSqlCacheManager().tryAddFeSqlCache(
- * statementContext.getConnectContext(),
- * statementContext.getOriginStatement().originStmt
- * );
- * }
- * </pre>
+ * The planner centrally handles SQL cache admission after an implementation
returns a result set.
*/
public interface ComputeResultSet {
- Optional<ResultSet> computeResultInFe(CascadesContext cascadesContext,
Optional<SqlCacheContext> sqlCacheContext,
- List<Slot> outputSlots);
+ Optional<ResultSet> computeResultInFe(CascadesContext cascadesContext,
List<Slot> outputSlots);
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalBlackholeSink.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalBlackholeSink.java
index de9e4190afb..0d2ddaad1e5 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalBlackholeSink.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalBlackholeSink.java
@@ -18,7 +18,6 @@
package org.apache.doris.nereids.trees.plans.physical;
import org.apache.doris.nereids.CascadesContext;
-import org.apache.doris.nereids.SqlCacheContext;
import org.apache.doris.nereids.memo.GroupExpression;
import org.apache.doris.nereids.properties.LogicalProperties;
import org.apache.doris.nereids.properties.PhysicalProperties;
@@ -133,11 +132,10 @@ public class PhysicalBlackholeSink<CHILD_TYPE extends
Plan> extends PhysicalSink
}
@Override
- public Optional<ResultSet> computeResultInFe(
- CascadesContext cascadesContext, Optional<SqlCacheContext>
sqlCacheContext, List<Slot> outputSlots) {
+ public Optional<ResultSet> computeResultInFe(CascadesContext
cascadesContext, List<Slot> outputSlots) {
CHILD_TYPE child = child();
if (child instanceof ComputeResultSet) {
- return ((ComputeResultSet)
child).computeResultInFe(cascadesContext, sqlCacheContext, outputSlots);
+ return ((ComputeResultSet)
child).computeResultInFe(cascadesContext, outputSlots);
} else {
return Optional.empty();
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalEmptyRelation.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalEmptyRelation.java
index 39bb72e1a99..ce3bb89ded9 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalEmptyRelation.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalEmptyRelation.java
@@ -18,10 +18,7 @@
package org.apache.doris.nereids.trees.plans.physical;
import org.apache.doris.catalog.Column;
-import org.apache.doris.catalog.Env;
import org.apache.doris.nereids.CascadesContext;
-import org.apache.doris.nereids.SqlCacheContext;
-import org.apache.doris.nereids.StatementContext;
import org.apache.doris.nereids.memo.GroupExpression;
import org.apache.doris.nereids.properties.LogicalProperties;
import org.apache.doris.nereids.properties.PhysicalProperties;
@@ -39,7 +36,6 @@ import org.apache.doris.nereids.util.Utils;
import org.apache.doris.qe.CommonResultSet;
import org.apache.doris.qe.ResultSet;
import org.apache.doris.qe.ResultSetMetaData;
-import org.apache.doris.qe.cache.CacheAnalyzer;
import org.apache.doris.statistics.model.Statistics;
import com.google.common.collect.ImmutableList;
@@ -118,26 +114,14 @@ public class PhysicalEmptyRelation extends
PhysicalRelation
}
@Override
- public Optional<ResultSet> computeResultInFe(CascadesContext
cascadesContext,
- Optional<SqlCacheContext> sqlCacheContext, List<Slot> outputSlots)
{
+ public Optional<ResultSet> computeResultInFe(CascadesContext
cascadesContext, List<Slot> outputSlots) {
List<Column> columns = Lists.newArrayList();
for (NamedExpression output : outputSlots) {
columns.add(new Column(output.getName(),
output.getDataType().toCatalogDataType()));
}
- StatementContext statementContext =
cascadesContext.getStatementContext();
- boolean enableSqlCache
- =
CacheAnalyzer.canUseSqlCache(statementContext.getConnectContext().getSessionVariable());
-
ResultSetMetaData metadata = new
CommonResultSet.CommonResultSetMetaData(columns);
ResultSet resultSet = new CommonResultSet(metadata,
ImmutableList.of());
- if (sqlCacheContext.isPresent() && enableSqlCache) {
- sqlCacheContext.get().setResultSetInFe(resultSet);
- Env.getCurrentEnv().getSqlCacheManager().tryAddFeSqlCache(
- statementContext.getConnectContext(),
- statementContext.getOriginStatement().originStmt
- );
- }
return Optional.of(resultSet);
}
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalOneRowRelation.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalOneRowRelation.java
index 371305641ce..a68cbb9d957 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalOneRowRelation.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalOneRowRelation.java
@@ -21,10 +21,7 @@ import org.apache.doris.analysis.ExprToStringValueVisitor;
import org.apache.doris.analysis.LiteralExpr;
import org.apache.doris.analysis.StringValueContext;
import org.apache.doris.catalog.Column;
-import org.apache.doris.catalog.Env;
import org.apache.doris.nereids.CascadesContext;
-import org.apache.doris.nereids.SqlCacheContext;
-import org.apache.doris.nereids.StatementContext;
import org.apache.doris.nereids.memo.GroupExpression;
import org.apache.doris.nereids.properties.DataTrait;
import org.apache.doris.nereids.properties.LogicalProperties;
@@ -46,7 +43,6 @@ import org.apache.doris.qe.CommonResultSet;
import org.apache.doris.qe.ConnectContext;
import org.apache.doris.qe.ResultSet;
import org.apache.doris.qe.ResultSetMetaData;
-import org.apache.doris.qe.cache.CacheAnalyzer;
import org.apache.doris.statistics.model.Statistics;
import com.google.common.collect.ImmutableList;
@@ -163,7 +159,7 @@ public class PhysicalOneRowRelation extends
PhysicalRelation implements OneRowRe
@Override
public Optional<ResultSet> computeResultInFe(
- CascadesContext cascadesContext, Optional<SqlCacheContext>
sqlCacheContext, List<Slot> outputSlots) {
+ CascadesContext cascadesContext, List<Slot> outputSlots) {
List<Column> columns = Lists.newArrayList();
List<String> data = Lists.newArrayList();
for (Slot outputSlot : outputSlots) {
@@ -201,16 +197,6 @@ public class PhysicalOneRowRelation extends
PhysicalRelation implements OneRowRe
ResultSetMetaData metadata = new
CommonResultSet.CommonResultSetMetaData(columns);
ResultSet resultSet = new CommonResultSet(metadata,
Collections.singletonList(data));
- StatementContext statementContext =
cascadesContext.getStatementContext();
- boolean enableSqlCache
- =
CacheAnalyzer.canUseSqlCache(statementContext.getConnectContext().getSessionVariable());
- if (sqlCacheContext.isPresent() && enableSqlCache) {
- sqlCacheContext.get().setResultSetInFe(resultSet);
- Env.getCurrentEnv().getSqlCacheManager().tryAddFeSqlCache(
- statementContext.getConnectContext(),
- statementContext.getOriginStatement().originStmt
- );
- }
return Optional.of(resultSet);
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalResultSink.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalResultSink.java
index 0c71298bbc2..a6e5e30b7f2 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalResultSink.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalResultSink.java
@@ -18,7 +18,6 @@
package org.apache.doris.nereids.trees.plans.physical;
import org.apache.doris.nereids.CascadesContext;
-import org.apache.doris.nereids.SqlCacheContext;
import org.apache.doris.nereids.memo.GroupExpression;
import org.apache.doris.nereids.properties.LogicalProperties;
import org.apache.doris.nereids.properties.PhysicalProperties;
@@ -130,11 +129,10 @@ public class PhysicalResultSink<CHILD_TYPE extends Plan>
extends PhysicalSink<CH
}
@Override
- public Optional<ResultSet> computeResultInFe(
- CascadesContext cascadesContext, Optional<SqlCacheContext>
sqlCacheContext, List<Slot> outputSlots) {
+ public Optional<ResultSet> computeResultInFe(CascadesContext
cascadesContext, List<Slot> outputSlots) {
CHILD_TYPE child = child();
if (child instanceof ComputeResultSet) {
- return ((ComputeResultSet)
child).computeResultInFe(cascadesContext, sqlCacheContext, outputSlots);
+ return ((ComputeResultSet)
child).computeResultInFe(cascadesContext, outputSlots);
} else {
return Optional.empty();
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalSqlCache.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalSqlCache.java
index 68075add788..1dc05e9cb8c 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalSqlCache.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/physical/PhysicalSqlCache.java
@@ -21,7 +21,6 @@ import org.apache.doris.analysis.Expr;
import org.apache.doris.common.util.DebugUtil;
import org.apache.doris.mysql.FieldInfo;
import org.apache.doris.nereids.CascadesContext;
-import org.apache.doris.nereids.SqlCacheContext;
import org.apache.doris.nereids.memo.GroupExpression;
import org.apache.doris.nereids.properties.DataTrait;
import org.apache.doris.nereids.properties.LogicalProperties;
@@ -168,8 +167,7 @@ public class PhysicalSqlCache extends PhysicalLeaf
}
@Override
- public Optional<ResultSet> computeResultInFe(
- CascadesContext cascadesContext, Optional<SqlCacheContext>
sqlCacheContext, List<Slot> outputSlots) {
+ public Optional<ResultSet> computeResultInFe(CascadesContext
cascadesContext, List<Slot> outputSlots) {
return resultSet;
}
}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/nereids/cost/CostModelV1Test.java
b/fe/fe-core/src/test/java/org/apache/doris/nereids/cost/CostModelV1Test.java
index 9ecf6cab8ae..8853c52b68e 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/nereids/cost/CostModelV1Test.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/nereids/cost/CostModelV1Test.java
@@ -18,6 +18,7 @@
package org.apache.doris.nereids.cost;
import org.apache.doris.nereids.PlanContext;
+import org.apache.doris.nereids.StatementContext;
import org.apache.doris.nereids.sqltest.SqlTestBase;
import org.apache.doris.nereids.trees.expressions.Slot;
import org.apache.doris.nereids.trees.expressions.functions.agg.AggregateParam;
@@ -28,6 +29,7 @@ import
org.apache.doris.nereids.trees.plans.physical.PhysicalHashAggregate;
import org.apache.doris.nereids.trees.plans.physical.PhysicalHashJoin;
import org.apache.doris.nereids.util.PlanChecker;
import org.apache.doris.nereids.util.PlanConstructor;
+import org.apache.doris.qe.SessionVariable;
import org.apache.doris.statistics.model.Statistics;
import org.apache.doris.statistics.model.StatisticsBuilder;
@@ -40,6 +42,51 @@ import java.util.Optional;
class CostModelV1Test extends SqlTestBase {
+ @Test
+ void testAddCostRecomputesWeightedValueFromComponents() {
+ CostWeight costWeight = new CostWeight(1, 1, 1.5, 1);
+ Cost planCost = Cost.of(costWeight, 0.1, 0.1, 0.1);
+ Cost childCost = Cost.of(costWeight, 0.1, 0.3, 0.7);
+
+ Cost totalCost = planCost.add(childCost, costWeight);
+ Cost expectedCost = Cost.of(costWeight,
+ planCost.getCpuCost() + childCost.getCpuCost(),
+ planCost.getMemoryCost() + childCost.getMemoryCost(),
+ planCost.getNetworkCost() + childCost.getNetworkCost());
+
+ Assertions.assertEquals(expectedCost.getValue(), totalCost.getValue());
+ Assertions.assertNotEquals(planCost.getValue() + childCost.getValue(),
totalCost.getValue());
+ Assertions.assertEquals(expectedCost.getCpuCost(),
totalCost.getCpuCost());
+ Assertions.assertEquals(expectedCost.getMemoryCost(),
totalCost.getMemoryCost());
+ Assertions.assertEquals(expectedCost.getNetworkCost(),
totalCost.getNetworkCost());
+ }
+
+ @Test
+ void testShareCostWeightInStatementContext() {
+ SessionVariable sessionVariable = connectContext.getSessionVariable();
+ double originalCpuWeight = sessionVariable.getCboCpuWeight();
+ StatementContext statementContext = new
StatementContext(connectContext, null);
+ try {
+ sessionVariable.setCboCpuWeight(2);
+ CostWeight firstWeight = statementContext.getCostWeight();
+
+ Assertions.assertSame(firstWeight,
statementContext.getCostWeight());
+ Assertions.assertEquals(2, Cost.ofCpu(firstWeight, 1).getValue());
+
+ sessionVariable.setCboCpuWeight(3);
+ Assertions.assertSame(firstWeight,
statementContext.getCostWeight());
+ Assertions.assertEquals(2, Cost.ofCpu(firstWeight, 1).getValue());
+
+ statementContext.setConnectContext(connectContext);
+ CostWeight nextExecutionWeight = statementContext.getCostWeight();
+
+ Assertions.assertNotSame(firstWeight, nextExecutionWeight);
+ Assertions.assertEquals(3, Cost.ofCpu(nextExecutionWeight,
1).getValue());
+ } finally {
+ sessionVariable.setCboCpuWeight(originalCpuWeight);
+ }
+ }
+
@Test
void testMaterializingCost() {
String sql = "select T1.id, T2.id, T2.score from T1 left join T2 "
@@ -74,6 +121,7 @@ class CostModelV1Test extends SqlTestBase {
PlanContext context = Mockito.mock(PlanContext.class);
Mockito.when(context.getChildStatistics(0)).thenReturn(childStats);
Mockito.when(context.getSessionVariable()).thenReturn(connectContext.getSessionVariable());
+
Mockito.when(context.getCostWeight()).thenReturn(connectContext.getStatementContext().getCostWeight());
Cost cost = new
CostModel(connectContext).visitPhysicalHashAggregate(aggregate, context);
Cost singlePointCost = new
CostModel(connectContext).visitPhysicalHashAggregate(singlePointAggregate,
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/nereids/minidump/MinidumpUtTest.java
b/fe/fe-core/src/test/java/org/apache/doris/nereids/minidump/MinidumpUtTest.java
index 7c01338955c..81fb1cd1bdf 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/nereids/minidump/MinidumpUtTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/nereids/minidump/MinidumpUtTest.java
@@ -20,11 +20,11 @@ package org.apache.doris.nereids.minidump;
import org.apache.doris.catalog.Env;
import org.apache.doris.common.Config;
import org.apache.doris.common.proc.FrontendsProcNode;
+import org.apache.doris.qe.ConnectContext;
import org.apache.doris.system.Frontend;
import org.json.JSONObject;
import org.junit.jupiter.api.Assertions;
-import org.junit.jupiter.api.Disabled;
import org.junit.jupiter.api.Test;
import org.mockito.MockedStatic;
import org.mockito.Mockito;
@@ -37,22 +37,28 @@ import org.mockito.Mockito;
*/
class MinidumpUtTest {
- @Disabled
@Test
public void testMinidumpUt() {
Minidump minidump = null;
String filePath =
getClass().getProtectionDomain().getCodeSource().getLocation().getPath();
String directory = filePath.substring(0,
filePath.indexOf("/target/test-classes"));
String currentMinidumpPath =
"/src/test/java/org/apache/doris/nereids/minidump/MinidumpUtTestData.json";
- try {
+ Frontend fe = Mockito.mock(Frontend.class);
+ Mockito.when(fe.getVersion()).thenReturn("Apache Doris test");
+ Env mockEnv = Mockito.mock(Env.class);
+ try (MockedStatic<Env> envStatic = Mockito.mockStatic(Env.class);
+ MockedStatic<FrontendsProcNode> procStatic =
Mockito.mockStatic(FrontendsProcNode.class)) {
+ envStatic.when(Env::getCurrentEnv).thenReturn(mockEnv);
+ procStatic.when(() ->
FrontendsProcNode.getCurrentFrontendVersion(mockEnv)).thenReturn(fe);
minidump = MinidumpUtils.jsonMinidumpLoad(directory +
currentMinidumpPath);
} catch (Exception e) {
throw new RuntimeException(e);
}
+ Assertions.assertNotNull(minidump);
MinidumpUtils.setConnectContext(minidump);
- JSONObject resultPlan = MinidumpUtils.executeSql("select * from t1
where l1 = 1");
- assert (minidump != null);
- assert (resultPlan != null);
+ Assertions.assertNull(ConnectContext.get().getStatementContext());
+ JSONObject resultPlan = MinidumpUtils.executeSql("select 1");
+ Assertions.assertNotNull(resultPlan);
}
@Test
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/nereids/minidump/MinidumpUtTestData.json
b/fe/fe-core/src/test/java/org/apache/doris/nereids/minidump/MinidumpUtTestData.json
index ed28fed7998..c7324437557 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/nereids/minidump/MinidumpUtTestData.json
+++
b/fe/fe-core/src/test/java/org/apache/doris/nereids/minidump/MinidumpUtTestData.json
@@ -1,4 +1,5 @@
{
+ "FeVersion": "test",
"Sql": "select * from t1",
"SessionVariable": {
"enable_nereids_planner": true,
@@ -10,7 +11,7 @@
"Tables": [
{
"TableType": "OLAP",
- "TableName": "t1",
+ "TableName": "[\"db\",\"t1\"]",
"TableValue": {
"clazz": "OlapTable",
"state": "NORMAL",
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/nereids/properties/ChildrenPropertiesRegulatorTest.java
b/fe/fe-core/src/test/java/org/apache/doris/nereids/properties/ChildrenPropertiesRegulatorTest.java
index 2bd445e60c6..66325eee97c 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/nereids/properties/ChildrenPropertiesRegulatorTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/nereids/properties/ChildrenPropertiesRegulatorTest.java
@@ -19,6 +19,7 @@ package org.apache.doris.nereids.properties;
import org.apache.doris.common.Pair;
import org.apache.doris.nereids.CascadesContext;
+import org.apache.doris.nereids.StatementContext;
import org.apache.doris.nereids.cost.Cost;
import org.apache.doris.nereids.cost.CostCalculator;
import org.apache.doris.nereids.jobs.JobContext;
@@ -63,7 +64,13 @@ public class ChildrenPropertiesRegulatorTest {
@BeforeEach
public void setUp() {
mockedJobContext = Mockito.mock(JobContext.class);
-
Mockito.when(mockedJobContext.getCascadesContext()).thenReturn(Mockito.mock(CascadesContext.class));
+ ConnectContext connectContext = new ConnectContext();
+ StatementContext statementContext = new
StatementContext(connectContext, null);
+ connectContext.setStatementContext(statementContext);
+ CascadesContext cascadesContext = Mockito.mock(CascadesContext.class);
+
Mockito.when(cascadesContext.getConnectContext()).thenReturn(connectContext);
+
Mockito.when(cascadesContext.getStatementContext()).thenReturn(statementContext);
+
Mockito.when(mockedJobContext.getCascadesContext()).thenReturn(cascadesContext);
}
@Test
@@ -91,10 +98,7 @@ public class ChildrenPropertiesRegulatorTest {
boolean canMergeChildProject) {
try (MockedStatic<CostCalculator> mockedCostCalculator =
Mockito.mockStatic(CostCalculator.class)) {
mockedCostCalculator.when(() ->
CostCalculator.calculateCost(Mockito.any(), Mockito.any(),
- Mockito.anyList())).thenReturn(Cost.zero());
- mockedCostCalculator.when(() ->
CostCalculator.addChildCost(Mockito.any(), Mockito.any(), Mockito.any(),
- Mockito.any(), Mockito.anyInt())).thenReturn(Cost.zero());
-
+ Mockito.anyList(), Mockito.any())).thenReturn(Cost.zero());
// project, cannot merge
Plan mockedChild = Mockito.mock(childClazz);
Mockito.when(mockedChild.withGroupExpression(Mockito.any())).thenReturn(mockedChild);
@@ -180,10 +184,7 @@ public class ChildrenPropertiesRegulatorTest {
private void testMustShuffleFilter(Class<? extends Plan> childClazz) {
try (MockedStatic<CostCalculator> mockedCostCalculator =
Mockito.mockStatic(CostCalculator.class)) {
mockedCostCalculator.when(() ->
CostCalculator.calculateCost(Mockito.any(), Mockito.any(),
- Mockito.anyList())).thenReturn(Cost.zero());
- mockedCostCalculator.when(() ->
CostCalculator.addChildCost(Mockito.any(), Mockito.any(), Mockito.any(),
- Mockito.any(), Mockito.anyInt())).thenReturn(Cost.zero());
-
+ Mockito.anyList(), Mockito.any())).thenReturn(Cost.zero());
// project, cannot merge
Plan mockedChild = Mockito.mock(childClazz);
Mockito.when(mockedChild.withGroupExpression(Mockito.any())).thenReturn(mockedChild);
diff --git a/fe/fe-core/src/test/java/org/apache/doris/qe/SqlCacheTest.java
b/fe/fe-core/src/test/java/org/apache/doris/qe/SqlCacheTest.java
index 087f671b3da..10d8b595f54 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/qe/SqlCacheTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/qe/SqlCacheTest.java
@@ -71,6 +71,10 @@ public class SqlCacheTest extends TestWithFeService {
Env currentEnv = Env.getCurrentEnv();
NereidsSqlCacheManager sqlCacheManager =
currentEnv.getSqlCacheManager();
Assertions.assertEquals(2,
sqlCacheManager.getSqlCaches().asMap().size());
+
sqlCacheManager.getSqlCaches().asMap().values().forEach(sqlCacheContext -> {
+ Assertions.assertNotNull(sqlCacheContext.getPhysicalPlan());
+
Assertions.assertFalse(sqlCacheContext.getPhysicalPlan().isEmpty());
+ });
executeNereidsSql("admin set frontend config
('sql_cache_manage_num'='1')");
Assertions.assertEquals(1,
sqlCacheManager.getSqlCaches().asMap().size());
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]