This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.2
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.2 by this push:
new be84e2afe85 branch:4.2: [improvement](planner) Reduce planner overhead
#67797 (#68290)
be84e2afe85 is described below
commit be84e2afe85ba44ffa5aade822f62e54b7c558c4
Author: morrySnow <[email protected]>
AuthorDate: Mon Sep 21 10:37:01 2026 +0800
branch:4.2: [improvement](planner) Reduce planner overhead #67797 (#68290)
picked from #67797
---
.../org/apache/doris/nereids/CascadesContext.java | 2 +-
.../org/apache/doris/nereids/NereidsPlanner.java | 22 ++++++--
.../java/org/apache/doris/nereids/PlanContext.java | 29 ++++++++++
.../org/apache/doris/nereids/StatementContext.java | 14 +++++
.../java/org/apache/doris/nereids/cost/Cost.java | 21 +++++---
.../apache/doris/nereids/cost/CostCalculator.java | 19 ++-----
.../org/apache/doris/nereids/cost/CostModel.java | 62 ++++++++++------------
.../nereids/jobs/cascades/CostAndEnforcerJob.java | 40 ++++----------
.../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 | 23 ++++----
.../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 | 10 ++--
.../doris/nereids/minidump/MinidumpUtTestData.json | 3 +-
.../ChildrenPropertiesRegulatorTest.java | 20 +++----
.../java/org/apache/doris/qe/SqlCacheTest.java | 4 ++
25 files changed, 247 insertions(+), 203 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 c9ef6a38d5a..0215f007022 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 cbbb918fa34..4fe629692ba 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
@@ -85,6 +85,7 @@ import org.apache.doris.qe.ConnectContext;
import org.apache.doris.qe.ResultSet;
import org.apache.doris.qe.SessionVariable;
import org.apache.doris.qe.VariableMgr;
+import org.apache.doris.qe.cache.CacheAnalyzer;
import org.apache.doris.statistics.util.StatisticsUtil;
import org.apache.doris.thrift.TQueryCacheParam;
@@ -740,7 +741,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;
}
@@ -1008,10 +1009,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;
}
}
@@ -1019,6 +1020,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 1f5dd91882c..7dc39be5923 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,7 +17,10 @@
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;
import org.apache.doris.nereids.trees.plans.Plan;
import org.apache.doris.qe.ConnectContext;
import org.apache.doris.qe.SessionVariable;
@@ -39,12 +42,18 @@ public class PlanContext {
private final int arity;
private boolean isBroadcastJoin = false;
private final boolean isStatsReliable;
+ private final CostWeight costWeight;
/**
* Constructor for PlanContext.
*/
public PlanContext(ConnectContext connectContext, GroupExpression
groupExpression) {
+ this(connectContext, groupExpression, (CostWeight) null);
+ }
+
+ private PlanContext(ConnectContext connectContext, GroupExpression
groupExpression, CostWeight costWeight) {
this.connectContext = connectContext;
+ this.costWeight = costWeight;
this.arity = groupExpression.arity();
this.planStats = groupExpression.getOwnerGroup().getStatistics();
this.isStatsReliable =
groupExpression.getOwnerGroup().isStatsReliable();
@@ -54,6 +63,21 @@ public class 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, groupExpression, costWeight);
+ if (childrenProperties.size() >= 2
+ && childrenProperties.get(1).getDistributionSpec() instanceof
DistributionSpecReplicated) {
+ isBroadcastJoin = true;
+ }
+ }
+
/**
* Constructor for plan-only context usage.
*/
@@ -63,6 +87,7 @@ public class PlanContext {
this.planStats = null;
this.arity = 0;
this.isStatsReliable = true;
+ this.costWeight = null;
}
public SessionVariable getSessionVariable() {
@@ -73,6 +98,10 @@ public class PlanContext {
isBroadcastJoin = true;
}
+ 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 dcb6d15060e..69a5803478a 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
@@ -39,6 +39,7 @@ import org.apache.doris.datasource.mvcc.MvccTableInfo;
import org.apache.doris.foundation.format.FormatOptions;
import org.apache.doris.mtmv.BaseTableInfo;
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;
@@ -130,6 +131,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 final Stopwatch stopwatch = Stopwatch.createUnstarted();
private final Stopwatch materializedViewStopwatch =
Stopwatch.createUnstarted();
@@ -520,6 +523,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) {
@@ -534,6 +540,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 Set<String> getUsedAIResourceNames() {
return Collections.unmodifiableSet(usedAIResourceNames);
}
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 7661e1aec51..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
@@ -19,11 +19,8 @@ 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.DistributionSpecReplicated;
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;
@@ -37,20 +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);
- if (childrenProperties.size() >= 2
- && childrenProperties.get(1).getDistributionSpec() instanceof
DistributionSpecReplicated) {
- planContext.setBroadcastJoin();
- }
-
+ 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 1e6393f162c..15b8acc7fec 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
@@ -104,14 +104,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();
@@ -134,7 +126,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())
@@ -148,10 +140,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) {
@@ -201,13 +193,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
@@ -215,14 +207,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
@@ -237,25 +229,25 @@ 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 visitPhysicalJdbcScan(PhysicalJdbcScan physicalJdbcScan,
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
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
public Cost visitPhysicalEsScan(PhysicalEsScan physicalEsScan, 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
@@ -271,7 +263,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
@@ -286,14 +278,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());
@@ -310,7 +302,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
@@ -322,7 +314,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);
@@ -331,7 +323,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);
@@ -339,7 +331,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
@@ -363,13 +355,13 @@ 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 {
int factor = aggregate.getGroupByExpressions().isEmpty() ? 1 :
beNumber;
// global
- return Cost.of(context.getSessionVariable(),
+ return Cost.of(context.getCostWeight(),
exprCost / 100 + inputStatistics.getRowCount() / factor,
inputStatistics.getRowCount() / factor, 0);
}
@@ -433,7 +425,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
);
@@ -494,14 +486,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
);
@@ -561,7 +553,7 @@ class CostModel extends PlanVisitor<Cost, PlanContext> {
if (leftStatistics.getRowCount() < 10 * rightStatistics.getRowCount())
{
nljPenalty = Math.min(leftStatistics.getRowCount(),
rightStatistics.getRowCount());
}
- return Cost.of(context.getSessionVariable(),
+ return Cost.of(context.getCostWeight(),
leftStatistics.getRowCount() * rightStatistics.getRowCount(),
rightStatistics.getRowCount() * nljPenalty,
0);
@@ -570,7 +562,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
@@ -585,13 +577,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
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 76a659870d0..0487ad74ae8 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,7 +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>
@@ -142,14 +140,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);
@@ -190,12 +180,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:
@@ -247,6 +231,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
@@ -267,17 +252,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 },
@@ -321,8 +303,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();
@@ -360,7 +343,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 b1cdfa49dcf..75802139b05 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.metrics.EventChannel;
import org.apache.doris.nereids.metrics.EventProducer;
import org.apache.doris.nereids.metrics.consumer.LogConsumer;
@@ -76,6 +77,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<>();
@@ -93,11 +95,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() {
@@ -942,15 +951,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 1602b79e54d..a065a85e982 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
@@ -220,10 +220,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);
@@ -241,8 +243,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 43aa2d92830..131f6e0026a 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;
@@ -815,10 +816,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 6dc6fe565a2..807c2580532 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() {
@@ -160,16 +163,14 @@ public class EnforceMissingPropertiesHelper {
oldOutputProperty, newOutputProperty);
ENFORCER_TRACER.log(EnforcerEvent.of(groupExpression, ((PhysicalPlan)
enforcer.getPlan()),
oldOutputProperty, newOutputProperty));
-
enforcer.setEstOutputRowCount(enforcer.getOwnerGroup().getStatistics().getRowCount());
- Cost enforcerCost = CostCalculator.calculateCost(connectContext,
enforcer,
- Lists.newArrayList(oldOutputProperty));
- enforcer.setCost(enforcerCost);
- curTotalCost = CostCalculator.addChildCost(
- connectContext,
- enforcer.getPlan(),
- enforcerCost,
- curTotalCost,
- 0);
+ Cost enforcerCost = enforcer.getCost();
+ if (enforcerCost == null) {
+
enforcer.setEstOutputRowCount(enforcer.getOwnerGroup().getStatistics().getRowCount());
+ enforcerCost = CostCalculator.calculateCost(connectContext,
enforcer,
+ Lists.newArrayList(oldOutputProperty), costWeight);
+ enforcer.setCost(enforcerCost);
+ }
+ 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/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 53666ccab14..02c0f000495 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;
@@ -132,11 +131,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 b1b5ea2beee..10d3788b9dc 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;
@@ -38,7 +35,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.Statistics;
import com.google.common.collect.ImmutableList;
@@ -117,26 +113,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 f61c5584625..d2c34af0ab1 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
@@ -19,10 +19,7 @@ package org.apache.doris.nereids.trees.plans.physical;
import org.apache.doris.analysis.LiteralExpr;
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;
@@ -43,7 +40,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.Statistics;
import com.google.common.collect.ImmutableList;
@@ -160,7 +156,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) {
@@ -197,16 +193,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 46df134c0cd..3c69e3efecb 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;
@@ -129,11 +128,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 afba3add8d7..1ec4ebb7807 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 9ebd933d46c..fd4785ecae0 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
@@ -17,15 +17,63 @@
package org.apache.doris.nereids.cost;
+import org.apache.doris.nereids.StatementContext;
import org.apache.doris.nereids.sqltest.SqlTestBase;
import org.apache.doris.nereids.trees.plans.Plan;
import org.apache.doris.nereids.trees.plans.physical.PhysicalHashJoin;
import org.apache.doris.nereids.util.PlanChecker;
+import org.apache.doris.qe.SessionVariable;
+import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
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 "
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 60ff399fa0b..5ed150b1e36 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
@@ -17,7 +17,10 @@
package org.apache.doris.nereids.minidump;
+import org.apache.doris.qe.ConnectContext;
+
import org.json.JSONObject;
+import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Disabled;
import org.junit.jupiter.api.Test;
@@ -41,9 +44,10 @@ class MinidumpUtTest {
} 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);
}
}
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 f63cc5d8d45..510ccdc6290 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;
@@ -29,6 +30,7 @@ import org.apache.doris.nereids.trees.plans.Plan;
import org.apache.doris.nereids.trees.plans.physical.PhysicalFilter;
import org.apache.doris.nereids.trees.plans.physical.PhysicalLimit;
import org.apache.doris.nereids.trees.plans.physical.PhysicalProject;
+import org.apache.doris.qe.ConnectContext;
import com.google.common.collect.Lists;
import com.google.common.collect.Maps;
@@ -54,7 +56,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
@@ -82,10 +90,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);
@@ -141,10 +146,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]