Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,6 @@
import org.apache.doris.qe.OriginStatement;
import org.apache.doris.qe.VariableMgr;
import org.apache.doris.resource.Tag;
import org.apache.doris.resource.computegroup.ComputeGroup;
import org.apache.doris.rpc.RpcException;
import org.apache.doris.service.FrontendOptions;
import org.apache.doris.statistics.AnalysisInfo;
Expand All @@ -75,16 +74,12 @@
import org.apache.doris.statistics.OlapAnalysisTask;
import org.apache.doris.statistics.util.StatisticsUtil;
import org.apache.doris.system.Backend;
import org.apache.doris.system.BeSelectionPolicy;
import org.apache.doris.system.SystemInfoService;
import org.apache.doris.thrift.TColumn;
import org.apache.doris.thrift.TCompressionType;
import org.apache.doris.thrift.TEncryptionAlgorithm;
import org.apache.doris.thrift.TFetchOption;
import org.apache.doris.thrift.TInvertedIndexFileStorageFormat;
import org.apache.doris.thrift.TNodeInfo;
import org.apache.doris.thrift.TOlapTable;
import org.apache.doris.thrift.TPaloNodesInfo;
import org.apache.doris.thrift.TPatternType;
import org.apache.doris.thrift.TPrimitiveType;
import org.apache.doris.thrift.TSortType;
Expand All @@ -106,7 +101,6 @@
import com.google.gson.annotations.SerializedName;
import lombok.Getter;
import lombok.Setter;
import org.apache.commons.collections4.CollectionUtils;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;

Expand Down Expand Up @@ -3174,48 +3168,6 @@ public AutoIncrementGenerator getAutoIncrementGenerator() {
return autoIncrementGenerator;
}

/**
* generate two phase read fetch option from this olap table.
*
* @param selectedIndexId the index want to scan
*/
public TFetchOption generateTwoPhaseReadOption(long selectedIndexId) {
boolean useStoreRow = this.storeRowColumn()
&& CollectionUtils.isEmpty(getTableProperty().getCopiedRowStoreColumns());
TFetchOption fetchOption = new TFetchOption();
fetchOption.setFetchRowStore(useStoreRow);
fetchOption.setUseTwoPhaseFetch(true);

ConnectContext context = ConnectContext.get();
if (context == null) {
context = new ConnectContext();
}
BeSelectionPolicy policy = new BeSelectionPolicy.Builder()
.needQueryAvailable()
.setRequireAliveBe()
.build();

TPaloNodesInfo nodesInfo = new TPaloNodesInfo();
ComputeGroup computeGroup = context.getComputeGroupSafely();

if (ComputeGroup.INVALID_COMPUTE_GROUP.equals(computeGroup)) {
throw new RuntimeException(ComputeGroup.INVALID_COMPUTE_GROUP_ERR_MSG);
}

for (Backend backend : policy.getCandidateBackends(computeGroup.getBackendList())) {
nodesInfo.addToNodes(new TNodeInfo(backend.getId(), 0, backend.getHost(), backend.getBrpcPort()));
}

fetchOption.setNodesInfo(nodesInfo);

if (!useStoreRow) {
List<TColumn> columnsDesc = Lists.newArrayList();
getColumnDesc(selectedIndexId, columnsDesc, null, null);
fetchOption.setColumnDesc(columnsDesc);
}
return fetchOption;
}

public void getColumnDesc(long selectedIndexId, List<TColumn> columnsDesc, List<String> keyColumnNames,
List<TPrimitiveType> keyColumnTypes, Set<String> materializedColumnNames) {
if (selectedIndexId != -1) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,8 +41,6 @@
import org.apache.doris.nereids.trees.plans.PlanNodeAndHash;
import org.apache.doris.nereids.trees.plans.algebra.OlapScan;
import org.apache.doris.nereids.trees.plans.physical.PhysicalAssertNumRows;
import org.apache.doris.nereids.trees.plans.physical.PhysicalDeferMaterializeOlapScan;
import org.apache.doris.nereids.trees.plans.physical.PhysicalDeferMaterializeTopN;
import org.apache.doris.nereids.trees.plans.physical.PhysicalDistribute;
import org.apache.doris.nereids.trees.plans.physical.PhysicalEsScan;
import org.apache.doris.nereids.trees.plans.physical.PhysicalFileScan;
Expand Down Expand Up @@ -207,12 +205,6 @@ public Cost visitPhysicalFilter(PhysicalFilter<? extends Plan> filter, PlanConte
(filter.getConjuncts().size() - prefixIndexMatched + exprCost) * filterCostFactor);
}

@Override
public Cost visitPhysicalDeferMaterializeOlapScan(PhysicalDeferMaterializeOlapScan deferMaterializeOlapScan,
PlanContext context) {
return visitPhysicalOlapScan(deferMaterializeOlapScan.getPhysicalOlapScan(), context);
}

public Cost visitPhysicalSchemaScan(PhysicalSchemaScan physicalSchemaScan, PlanContext context) {
Statistics statistics = context.getStatisticsWithCheck();
return Cost.ofCpu(context.getSessionVariable(), statistics.getRowCount());
Expand Down Expand Up @@ -297,12 +289,6 @@ public Cost visitPhysicalTopN(PhysicalTopN<? extends Plan> topN, PlanContext con
return Cost.of(context.getSessionVariable(), childRowCount, rowCount, childRowCount);
}

@Override
public Cost visitPhysicalDeferMaterializeTopN(PhysicalDeferMaterializeTopN<? extends Plan> topN,
PlanContext context) {
return visitPhysicalTopN(topN.getPhysicalTopN(), context);
}

@Override
public Cost visitPhysicalPartitionTopN(PhysicalPartitionTopN<? extends Plan> partitionTopN, PlanContext context) {
Statistics statistics = context.getStatisticsWithCheck();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -133,9 +133,6 @@
import org.apache.doris.nereids.trees.plans.physical.PhysicalCTEAnchor;
import org.apache.doris.nereids.trees.plans.physical.PhysicalCTEConsumer;
import org.apache.doris.nereids.trees.plans.physical.PhysicalCTEProducer;
import org.apache.doris.nereids.trees.plans.physical.PhysicalDeferMaterializeOlapScan;
import org.apache.doris.nereids.trees.plans.physical.PhysicalDeferMaterializeResultSink;
import org.apache.doris.nereids.trees.plans.physical.PhysicalDeferMaterializeTopN;
import org.apache.doris.nereids.trees.plans.physical.PhysicalDictionarySink;
import org.apache.doris.nereids.trees.plans.physical.PhysicalDistribute;
import org.apache.doris.nereids.trees.plans.physical.PhysicalEmptyRelation;
Expand Down Expand Up @@ -248,7 +245,6 @@
import org.apache.doris.tablefunction.TableValuedFunctionIf;
import org.apache.doris.thrift.TExternalTableSinkHashAlgorithm;
import org.apache.doris.thrift.TExternalTableSinkWriterAssignment;
import org.apache.doris.thrift.TFetchOption;
import org.apache.doris.thrift.TPaimonFixedBucketInfo;
import org.apache.doris.thrift.TPartitionType;
import org.apache.doris.thrift.TPushAggOp;
Expand Down Expand Up @@ -508,16 +504,6 @@ public PlanFragment visitPhysicalResultSink(PhysicalResultSink<? extends Plan> p
return planFragment;
}

@Override
public PlanFragment visitPhysicalDeferMaterializeResultSink(
PhysicalDeferMaterializeResultSink<? extends Plan> sink,
PlanTranslatorContext context) {
PlanFragment planFragment = visitPhysicalResultSink(sink.getPhysicalResultSink(), context);
TFetchOption fetchOption = sink.getOlapTable().generateTwoPhaseReadOption(sink.getSelectedIndexId());
((ResultSink) planFragment.getSink()).setFetchOption(fetchOption);
return planFragment;
}

@Override
public PlanFragment visitPhysicalDictionarySink(PhysicalDictionarySink<? extends Plan> dictionarySink,
PlanTranslatorContext context) {
Expand Down Expand Up @@ -1146,17 +1132,6 @@ private void translateRuntimeFilter(PhysicalRelation physicalRelation, ScanNode
context.getTopnFilterContext().translateTarget(physicalRelation, scanNode, context);
}

@Override
public PlanFragment visitPhysicalDeferMaterializeOlapScan(
PhysicalDeferMaterializeOlapScan deferMaterializeOlapScan, PlanTranslatorContext context) {
PlanFragment planFragment = visitPhysicalOlapScan(deferMaterializeOlapScan.getPhysicalOlapScan(), context);
OlapScanNode olapScanNode = (OlapScanNode) planFragment.getPlanRoot();
TupleDescriptor tupleDescriptor = context.getTupleDesc(olapScanNode.getTupleId());
context.createSlotDesc(tupleDescriptor, deferMaterializeOlapScan.getColumnIdSlot());
context.getTopnFilterContext().translateTarget(deferMaterializeOlapScan, olapScanNode, context);
return planFragment;
}

@Override
public PlanFragment visitPhysicalOneRowRelation(PhysicalOneRowRelation oneRowRelation,
PlanTranslatorContext context) {
Expand Down Expand Up @@ -2316,12 +2291,7 @@ public PlanFragment visitPhysicalProject(PhysicalProject<? extends Plan> project
List<Slot> slots = null;
// TODO FE/BE do not support multi-layer-project on MultiDataSink now.
if (project.hasMultiLayerProjection()
&& !(inputFragment instanceof MultiCastPlanFragment)
// TODO support for two phase read with project, remove it after refactor
&& !(project.child() instanceof PhysicalDeferMaterializeTopN)
&& !(project.child() instanceof PhysicalDeferMaterializeOlapScan
|| (project.child() instanceof PhysicalFilter
&& ((PhysicalFilter<?>) project.child()).child() instanceof PhysicalDeferMaterializeOlapScan))) {
&& !(inputFragment instanceof MultiCastPlanFragment)) {
int layerCount = project.getMultiLayerProjects().size();
for (int i = 0; i < layerCount; i++) {
List<NamedExpression> layer = project.getMultiLayerProjects().get(i);
Expand Down Expand Up @@ -2438,28 +2408,20 @@ public PlanFragment visitPhysicalProject(PhysicalProject<? extends Plan> project
}

if (inputPlanNode instanceof ScanNode) {
// TODO support for two phase read with project, remove this if after refactor
if (!(project.child() instanceof PhysicalDeferMaterializeOlapScan
|| (project.child() instanceof PhysicalFilter
&& ((PhysicalFilter<?>) project.child()).child() instanceof PhysicalDeferMaterializeOlapScan))) {
TupleDescriptor projectionTuple = generateTupleDesc(slots,
((ScanNode) inputPlanNode).getTupleDesc().getTable(), context);
inputPlanNode.setProjectList(projectionExprs);
inputPlanNode.setOutputTupleDesc(projectionTuple);
}
TupleDescriptor projectionTuple = generateTupleDesc(slots,
((ScanNode) inputPlanNode).getTupleDesc().getTable(), context);
inputPlanNode.setProjectList(projectionExprs);
inputPlanNode.setOutputTupleDesc(projectionTuple);

if (inputPlanNode instanceof OlapScanNode) {
((OlapScanNode) inputPlanNode).updateRequiredSlots(context, requiredByProjectSlotIdSet);
}
updateScanSlotsMaterialization((ScanNode) inputPlanNode, requiredSlotIdSet,
requiredByProjectSlotIdSet, context);
} else {
if (project.child() instanceof PhysicalDeferMaterializeTopN) {
inputFragment.setOutputExprs(allProjectionExprs);
} else {
TupleDescriptor tupleDescriptor = generateTupleDesc(slots, null, context);
inputPlanNode.setProjectList(projectionExprs);
inputPlanNode.setOutputTupleDesc(tupleDescriptor);
}
TupleDescriptor tupleDescriptor = generateTupleDesc(slots, null, context);
inputPlanNode.setProjectList(projectionExprs);
inputPlanNode.setOutputTupleDesc(tupleDescriptor);
}
return inputFragment;
}
Expand Down Expand Up @@ -2750,21 +2712,6 @@ public PlanFragment visitPhysicalTopN(PhysicalTopN<? extends Plan> topN, PlanTra
return inputFragment;
}

@Override
public PlanFragment visitPhysicalDeferMaterializeTopN(PhysicalDeferMaterializeTopN<? extends Plan> topN,
PlanTranslatorContext context) {
PlanFragment planFragment = visitPhysicalTopN(topN.getPhysicalTopN(), context);
if (planFragment.getPlanRoot() instanceof SortNode) {
SortNode sortNode = (SortNode) planFragment.getPlanRoot();
sortNode.setUseTwoPhaseReadOpt(true);
sortNode.getSortInfo().setUseTwoPhaseRead();
if (context.getTopnFilterContext().isTopnFilterSource(topN)) {
context.getTopnFilterContext().translateSource(topN, sortNode);
}
}
return planFragment;
}

@Override
public PlanFragment visitPhysicalRepeat(PhysicalRepeat<? extends Plan> repeat, PlanTranslatorContext context) {
PlanFragment inputPlanFragment = repeat.child(0).accept(this, context);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,6 @@
import org.apache.doris.nereids.rules.rewrite.CreatePartitionTopNFromWindow;
import org.apache.doris.nereids.rules.rewrite.DecomposeRepeatWithPreAggregation;
import org.apache.doris.nereids.rules.rewrite.DecoupleEncodeDecode;
import org.apache.doris.nereids.rules.rewrite.DeferMaterializeTopNResult;
import org.apache.doris.nereids.rules.rewrite.DistinctAggStrategySelector;
import org.apache.doris.nereids.rules.rewrite.DistinctAggregateRewriter;
import org.apache.doris.nereids.rules.rewrite.DistinctWindowExpression;
Expand Down Expand Up @@ -797,9 +796,6 @@ public class Rewriter extends AbstractBatchJobExecutor {
topDown(new PushDownScoreTopNIntoOlapScan(),
new CheckScoreUsage())
),
topic("topn optimize",
topDown(new DeferMaterializeTopNResult())
),
topic("add projection for join",
custom(RuleType.ADD_PROJECT_FOR_JOIN, AddProjectForJoin::new),
topDown(new MergeProjectable())
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,6 @@
import org.apache.doris.nereids.trees.plans.logical.LogicalAggregate;
import org.apache.doris.nereids.trees.plans.logical.LogicalCTEConsumer;
import org.apache.doris.nereids.trees.plans.logical.LogicalCatalogRelation;
import org.apache.doris.nereids.trees.plans.logical.LogicalDeferMaterializeTopN;
import org.apache.doris.nereids.trees.plans.logical.LogicalFilter;
import org.apache.doris.nereids.trees.plans.logical.LogicalJoin;
import org.apache.doris.nereids.trees.plans.logical.LogicalSort;
Expand Down Expand Up @@ -352,21 +351,6 @@ public Void visitLogicalTopN(LogicalTopN<? extends Plan> topN, LineageInfo linea
return super.visitLogicalTopN(topN, lineageInfo);
}

/**
* Collect SORT indirect lineage from defer-materialize TopN.
*
* <p>Using the example SQL above, there is no defer-materialize TopN.
*/
@Override
public Void visitLogicalDeferMaterializeTopN(LogicalDeferMaterializeTopN<? extends Plan> topN,
LineageInfo lineageInfo) {
List<Expression> sortExprs = new ArrayList<>();
topN.getOrderKeys().forEach(orderKey -> sortExprs.add(orderKey.getExpr()));
Set<Expression> shuttled = shuttleExpressions(sortExprs, exprIdExpressionMap);
addIndirectLineage(IndirectLineageType.SORT, shuttled, lineageInfo);
return super.visitLogicalDeferMaterializeTopN(topN, lineageInfo);
}

/**
* Replace UNION outputs in direct lineage using child outputs.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,6 @@
import org.apache.doris.nereids.trees.plans.physical.PhysicalBlackholeSink;
import org.apache.doris.nereids.trees.plans.physical.PhysicalCTEAnchor;
import org.apache.doris.nereids.trees.plans.physical.PhysicalCTEProducer;
import org.apache.doris.nereids.trees.plans.physical.PhysicalDeferMaterializeResultSink;
import org.apache.doris.nereids.trees.plans.physical.PhysicalDictionarySink;
import org.apache.doris.nereids.trees.plans.physical.PhysicalDistribute;
import org.apache.doris.nereids.trees.plans.physical.PhysicalFilter;
Expand Down Expand Up @@ -303,13 +302,6 @@ public Plan visitPhysicalDictionarySink(PhysicalDictionarySink<? extends Plan> d
return rewriteUnary(dictionarySink, ctx.withAllowShuffleKeyPrune(childAllowShuffleKeyPrune));
}

@Override
public Plan visitPhysicalDeferMaterializeResultSink(
PhysicalDeferMaterializeResultSink<? extends Plan> sink,
PruneCtx ctx) {
return rewriteUnary(sink, ctx.withAllowShuffleKeyPrune(false));
}

private <P extends PhysicalUnary<?>> P rewriteUnary(P plan, PruneCtx ctx) {
Plan oldChild = plan.child();
Plan newChild = oldChild.accept(this, ctx);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,6 @@
import org.apache.doris.nereids.trees.plans.Plan;
import org.apache.doris.nereids.trees.plans.SortPhase;
import org.apache.doris.nereids.trees.plans.algebra.TopN;
import org.apache.doris.nereids.trees.plans.physical.PhysicalDeferMaterializeTopN;
import org.apache.doris.nereids.trees.plans.physical.PhysicalTopN;

/**
Expand All @@ -50,17 +49,11 @@ public PhysicalTopN<? extends Plan> visitPhysicalTopN(PhysicalTopN<? extends Pla
}

boolean checkTopN(TopN topN) {
if (!(topN instanceof PhysicalTopN) && !(topN instanceof PhysicalDeferMaterializeTopN)) {
if (!(topN instanceof PhysicalTopN)) {
return false;
}
if (topN instanceof PhysicalTopN
&& ((PhysicalTopN) topN).getSortPhase() != SortPhase.LOCAL_SORT) {
if (((PhysicalTopN) topN).getSortPhase() != SortPhase.LOCAL_SORT) {
return false;
} else {
if (topN instanceof PhysicalDeferMaterializeTopN
&& ((PhysicalDeferMaterializeTopN) topN).getSortPhase() != SortPhase.LOCAL_SORT) {
return false;
}
}

if (topN.getOrderKeys().isEmpty()) {
Expand All @@ -76,17 +69,4 @@ boolean checkTopN(TopN topN) {
return true;
}

@Override
public Plan visitPhysicalDeferMaterializeTopN(PhysicalDeferMaterializeTopN<? extends Plan> topN,
CascadesContext ctx) {
topN.child().accept(this, ctx);
if (checkTopN(topN)) {
TopnFilterPushDownVisitor pusher = new TopnFilterPushDownVisitor(ctx.getTopnFilterContext());
TopnFilterPushDownVisitor.PushDownContext pushdownContext = new PushDownContext(topN,
topN.getOrderKeys().get(0).getExpr(),
topN.getOrderKeys().get(0).isNullFirst());
topN.accept(pusher, pushdownContext);
}
return topN;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,6 @@
import org.apache.doris.nereids.trees.plans.algebra.Union;
import org.apache.doris.nereids.trees.plans.physical.PhysicalCTEAnchor;
import org.apache.doris.nereids.trees.plans.physical.PhysicalCTEProducer;
import org.apache.doris.nereids.trees.plans.physical.PhysicalDeferMaterializeOlapScan;
import org.apache.doris.nereids.trees.plans.physical.PhysicalEsScan;
import org.apache.doris.nereids.trees.plans.physical.PhysicalFileScan;
import org.apache.doris.nereids.trees.plans.physical.PhysicalHashJoin;
Expand Down Expand Up @@ -269,7 +268,6 @@ private boolean supportPhysicalRelations(PhysicalRelation relation) {
|| relation instanceof PhysicalEsScan
|| relation instanceof PhysicalFileScan
|| relation instanceof PhysicalJdbcScan
|| relation instanceof PhysicalDeferMaterializeOlapScan
|| relation instanceof PhysicalLazyMaterializeOlapScan;
}
}
Loading
Loading