Repository: calcite Updated Branches: refs/heads/master 6e8bb5a16 -> 68ba411e2
[CALCITE-2291] Support Push Project past Correlate (Chunhui Shi) Fix dynamic row type in UNNEST For the validation phase to pass use a place holder name for SqlUnnestOperator "$unnest" Use the same name as the input name for Uncollect Node Fixing the converter tests to expect "$unnest" (Hanumath Maduri) Fix ProjectCorrelateTransposeRule for SEMI and ANTI Correlate Fix bug with names loss for Unnest Update RelDataType for RexCorrelVariable when rule is applied (Volodymyr Vysotskyi) Close apache/calcite#685 Co-authored-by: chunhui-shi <[email protected]> Co-authored-by: HanumathRao <[email protected]> Co-authored-by: Volodymyr Vysotskyi <[email protected]> Project: http://git-wip-us.apache.org/repos/asf/calcite/repo Commit: http://git-wip-us.apache.org/repos/asf/calcite/commit/68ba411e Tree: http://git-wip-us.apache.org/repos/asf/calcite/tree/68ba411e Diff: http://git-wip-us.apache.org/repos/asf/calcite/diff/68ba411e Branch: refs/heads/master Commit: 68ba411e23ba930bb2086bb3eed4c46edfac23eb Parents: 6e8bb5a Author: chunhui-shi <[email protected]> Authored: Mon Apr 30 10:21:14 2018 -0700 Committer: Volodymyr Vysotskyi <[email protected]> Committed: Wed Jun 27 21:33:39 2018 +0300 ---------------------------------------------------------------------- .../org/apache/calcite/rel/core/Uncollect.java | 7 +- .../rules/ProjectCorrelateTransposeRule.java | 212 +++++++++++++++++++ .../apache/calcite/rel/rules/PushProjector.java | 45 +++- .../apache/calcite/sql/SqlUnnestOperator.java | 11 +- .../calcite/sql2rel/SqlToRelConverter.java | 2 +- .../apache/calcite/test/RelOptRulesTest.java | 153 +++++++++++++ .../org/apache/calcite/test/RelOptTestBase.java | 35 +++ .../calcite/test/SqlToRelConverterTest.java | 59 ++---- .../apache/calcite/test/SqlToRelTestBase.java | 47 ++++ .../org/apache/calcite/test/RelOptRulesTest.xml | 179 ++++++++++++++++ .../calcite/test/SqlToRelConverterTest.xml | 42 +++- 11 files changed, 731 insertions(+), 61 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/calcite/blob/68ba411e/core/src/main/java/org/apache/calcite/rel/core/Uncollect.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/calcite/rel/core/Uncollect.java b/core/src/main/java/org/apache/calcite/rel/core/Uncollect.java index 0c531cb..98fa388 100644 --- a/core/src/main/java/org/apache/calcite/rel/core/Uncollect.java +++ b/core/src/main/java/org/apache/calcite/rel/core/Uncollect.java @@ -23,7 +23,6 @@ import org.apache.calcite.rel.RelInput; import org.apache.calcite.rel.RelNode; import org.apache.calcite.rel.RelWriter; import org.apache.calcite.rel.SingleRel; -import org.apache.calcite.rel.type.DynamicRecordType; import org.apache.calcite.rel.type.RelDataType; import org.apache.calcite.rel.type.RelDataTypeFactory; import org.apache.calcite.rel.type.RelDataTypeField; @@ -129,10 +128,10 @@ public class Uncollect extends SingleRel { if (fields.size() == 1 && fields.get(0).getType().getSqlTypeName() == SqlTypeName.ANY) { - // Component type is unknown to Uncollect, build dynamic star record - // type. Only consider ONE field case for unknown type. + // Component type is unknown to Uncollect, build a row type with input column name + // and Any type. return builder - .add(DynamicRecordType.DYNAMIC_STAR_PREFIX, SqlTypeName.ANY) + .add(fields.get(0).getName(), SqlTypeName.ANY) .nullable(true) .build(); } http://git-wip-us.apache.org/repos/asf/calcite/blob/68ba411e/core/src/main/java/org/apache/calcite/rel/rules/ProjectCorrelateTransposeRule.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/calcite/rel/rules/ProjectCorrelateTransposeRule.java b/core/src/main/java/org/apache/calcite/rel/rules/ProjectCorrelateTransposeRule.java new file mode 100644 index 0000000..67d9402 --- /dev/null +++ b/core/src/main/java/org/apache/calcite/rel/rules/ProjectCorrelateTransposeRule.java @@ -0,0 +1,212 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to you under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.calcite.rel.rules; + +import org.apache.calcite.plan.RelOptRule; +import org.apache.calcite.plan.RelOptRuleCall; +import org.apache.calcite.plan.hep.HepRelVertex; +import org.apache.calcite.plan.volcano.RelSubset; +import org.apache.calcite.rel.RelNode; +import org.apache.calcite.rel.RelShuttleImpl; +import org.apache.calcite.rel.core.Correlate; +import org.apache.calcite.rel.core.CorrelationId; +import org.apache.calcite.rel.core.Project; +import org.apache.calcite.rel.core.RelFactories; +import org.apache.calcite.rex.RexBuilder; +import org.apache.calcite.rex.RexCorrelVariable; +import org.apache.calcite.rex.RexFieldAccess; +import org.apache.calcite.rex.RexNode; +import org.apache.calcite.rex.RexShuttle; +import org.apache.calcite.tools.RelBuilderFactory; +import org.apache.calcite.util.BitSets; +import org.apache.calcite.util.ImmutableBitSet; +import org.apache.calcite.util.Util; + +import java.util.BitSet; +import java.util.HashMap; +import java.util.Map; + +/** + * Push Project under Correlate to apply on Correlate's left and right child + */ +public class ProjectCorrelateTransposeRule extends RelOptRule { + + public static final ProjectCorrelateTransposeRule INSTANCE = + new ProjectCorrelateTransposeRule( + PushProjector.ExprCondition.TRUE, + RelFactories.LOGICAL_BUILDER); + + //~ Instance fields -------------------------------------------------------- + + /** + * preserveExprCondition to define the condition for a expression not to be pushed + */ + private final PushProjector.ExprCondition preserveExprCondition; + + //~ Constructors ----------------------------------------------------------- + + public ProjectCorrelateTransposeRule( + PushProjector.ExprCondition preserveExprCondition, + RelBuilderFactory relFactory) { + super( + operand(Project.class, + operand(Correlate.class, any())), + relFactory, null); + this.preserveExprCondition = preserveExprCondition; + } + + //~ Methods ---------------------------------------------------------------- + + public void onMatch(RelOptRuleCall call) { + Project origProj = call.rel(0); + final Correlate corr = call.rel(1); + + // locate all fields referenced in the projection + // determine which inputs are referenced in the projection; + // if all fields are being referenced and there are no + // special expressions, no point in proceeding any further + PushProjector pushProject = + new PushProjector( + origProj, + call.builder().literal(true), + corr, + preserveExprCondition, + call.builder()); + if (pushProject.locateAllRefs()) { + return; + } + + // create left and right projections, projecting only those + // fields referenced on each side + RelNode leftProjRel = + pushProject.createProjectRefsAndExprs( + corr.getLeft(), + true, + false); + RelNode rightProjRel = + pushProject.createProjectRefsAndExprs( + corr.getRight(), + true, + true); + + Map<Integer, Integer> requiredColsMap = new HashMap<>(); + + // adjust requiredColumns that reference the projected columns + int[] adjustments = pushProject.getAdjustments(); + BitSet updatedBits = new BitSet(); + for (Integer col : corr.getRequiredColumns()) { + int newCol = col + adjustments[col]; + updatedBits.set(newCol); + requiredColsMap.put(col, newCol); + } + + RexBuilder rexBuilder = call.builder().getRexBuilder(); + + CorrelationId correlationId = corr.getCluster().createCorrel(); + RexCorrelVariable rexCorrel = + (RexCorrelVariable) rexBuilder.makeCorrel( + leftProjRel.getRowType(), + correlationId); + + // updates RexCorrelVariable and sets actual RelDataType for RexFieldAccess + rightProjRel = rightProjRel.accept( + new RelNodesExprsHandler( + new RexFieldAccessReplacer(corr.getCorrelationId(), + rexCorrel, rexBuilder, requiredColsMap))); + + // create a new correlate with the projected children + Correlate newCorrRel = + corr.copy( + corr.getTraitSet(), + leftProjRel, + rightProjRel, + correlationId, + ImmutableBitSet.of(BitSets.toIter(updatedBits)), + corr.getJoinType()); + + // put the original project on top of the correlate, converting it to + // reference the modified projection list + RelNode topProject = + pushProject.createNewProject(newCorrRel, adjustments); + + call.transformTo(topProject); + } + + /** + * Visitor for RexNodes which replaces {@link RexCorrelVariable} with specified. + */ + public static class RexFieldAccessReplacer extends RexShuttle { + private final RexBuilder builder; + private final CorrelationId rexCorrelVariableToReplace; + private final RexCorrelVariable rexCorrelVariable; + private final Map<Integer, Integer> requiredColsMap; + + public RexFieldAccessReplacer( + CorrelationId rexCorrelVariableToReplace, + RexCorrelVariable rexCorrelVariable, + RexBuilder builder, + Map<Integer, Integer> requiredColsMap) { + this.rexCorrelVariableToReplace = rexCorrelVariableToReplace; + this.rexCorrelVariable = rexCorrelVariable; + this.builder = builder; + this.requiredColsMap = requiredColsMap; + } + + @Override public RexNode visitCorrelVariable(RexCorrelVariable variable) { + if (variable.id.equals(rexCorrelVariableToReplace)) { + return rexCorrelVariable; + } + return variable; + } + + @Override public RexNode visitFieldAccess(RexFieldAccess fieldAccess) { + RexNode refExpr = fieldAccess.getReferenceExpr().accept(this); + // creates new RexFieldAccess instance for the case when referenceExpr was replaced. + // Otherwise calls super method. + if (refExpr == rexCorrelVariable) { + return builder.makeFieldAccess( + refExpr, + requiredColsMap.get(fieldAccess.getField().getIndex())); + } + return super.visitFieldAccess(fieldAccess); + } + } + + /** + * Visitor for RelNodes which applies specified {@link RexShuttle} visitor + * for every node in the tree. + */ + public static class RelNodesExprsHandler extends RelShuttleImpl { + private final RexShuttle rexVisitor; + + public RelNodesExprsHandler(RexShuttle rexVisitor) { + this.rexVisitor = rexVisitor; + } + + @Override protected RelNode visitChild(RelNode parent, int i, RelNode child) { + if (child instanceof HepRelVertex) { + child = ((HepRelVertex) child).getCurrentRel(); + } else if (child instanceof RelSubset) { + RelSubset subset = (RelSubset) child; + child = Util.first(subset.getBest(), subset.getOriginal()); + } + return super.visitChild(parent, i, child).accept(rexVisitor); + } + } +} + +// End ProjectCorrelateTransposeRule.java http://git-wip-us.apache.org/repos/asf/calcite/blob/68ba411e/core/src/main/java/org/apache/calcite/rel/rules/PushProjector.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/calcite/rel/rules/PushProjector.java b/core/src/main/java/org/apache/calcite/rel/rules/PushProjector.java index 8eaf1a4..60a9663 100644 --- a/core/src/main/java/org/apache/calcite/rel/rules/PushProjector.java +++ b/core/src/main/java/org/apache/calcite/rel/rules/PushProjector.java @@ -20,6 +20,7 @@ import org.apache.calcite.linq4j.Ord; import org.apache.calcite.plan.RelOptUtil; import org.apache.calcite.plan.Strong; import org.apache.calcite.rel.RelNode; +import org.apache.calcite.rel.core.Correlate; import org.apache.calcite.rel.core.Join; import org.apache.calcite.rel.core.Project; import org.apache.calcite.rel.core.SemiJoin; @@ -32,6 +33,7 @@ import org.apache.calcite.rex.RexNode; import org.apache.calcite.rex.RexUtil; import org.apache.calcite.rex.RexVisitorImpl; import org.apache.calcite.runtime.PredicateImpl; +import org.apache.calcite.sql.SemiJoinType; import org.apache.calcite.sql.SqlOperator; import org.apache.calcite.tools.RelBuilder; import org.apache.calcite.util.BitSets; @@ -250,6 +252,45 @@ public class PushProjector { strongBitmap = ImmutableBitSet.range(nSysFields, nChildFields); } + } else if (childRel instanceof Correlate) { + Correlate corrRel = (Correlate) childRel; + List<RelDataTypeField> leftFields = + corrRel.getLeft().getRowType().getFieldList(); + List<RelDataTypeField> rightFields = + corrRel.getRight().getRowType().getFieldList(); + nFields = leftFields.size(); + SemiJoinType joinType = corrRel.getJoinType(); + switch (joinType) { + case SEMI: + case ANTI: + nFieldsRight = 0; + break; + default: + nFieldsRight = rightFields.size(); + } + nSysFields = 0; + childBitmap = + ImmutableBitSet.range(0, nFields); + rightBitmap = + ImmutableBitSet.range(nFields, nChildFields); + + // Required columns need to be included in project + projRefs.or(BitSets.of(corrRel.getRequiredColumns())); + + switch (joinType) { + case INNER: + strongBitmap = ImmutableBitSet.of(); + break; + case ANTI: + case SEMI: // All the left-input's columns must be strong + strongBitmap = ImmutableBitSet.range(0, nFields); + break; + case LEFT: // All the right-input's columns must be strong + strongBitmap = ImmutableBitSet.range(nFields, nChildFields); + break; + default: + strongBitmap = ImmutableBitSet.range(0, nChildFields); + } } else { nFields = nChildFields; nFieldsRight = 0; @@ -260,8 +301,8 @@ public class PushProjector { } assert nChildFields == nSysFields + nFields + nFieldsRight; - childPreserveExprs = new ArrayList<RexNode>(); - rightPreserveExprs = new ArrayList<RexNode>(); + childPreserveExprs = new ArrayList<>(); + rightPreserveExprs = new ArrayList<>(); rexBuilder = childRel.getCluster().getRexBuilder(); } http://git-wip-us.apache.org/repos/asf/calcite/blob/68ba411e/core/src/main/java/org/apache/calcite/sql/SqlUnnestOperator.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/calcite/sql/SqlUnnestOperator.java b/core/src/main/java/org/apache/calcite/sql/SqlUnnestOperator.java index 6825bc1..d55c07e 100644 --- a/core/src/main/java/org/apache/calcite/sql/SqlUnnestOperator.java +++ b/core/src/main/java/org/apache/calcite/sql/SqlUnnestOperator.java @@ -16,7 +16,6 @@ */ package org.apache.calcite.sql; -import org.apache.calcite.rel.type.DynamicRecordType; import org.apache.calcite.rel.type.RelDataType; import org.apache.calcite.rel.type.RelDataTypeFactory; import org.apache.calcite.sql.type.ArraySqlType; @@ -65,13 +64,13 @@ public class SqlUnnestOperator extends SqlFunctionalOperator { for (Integer operand : Util.range(opBinding.getOperandCount())) { RelDataType type = opBinding.getOperandType(operand); if (type.getSqlTypeName() == SqlTypeName.ANY) { - // When there is one operand with unknown type (ANY), the return type - // is dynamic star + // Unnest Operator in schema less systems returns one column as the output + // $unnest is a place holder to specify that one column with type ANY is output. return builder - .add(DynamicRecordType.DYNAMIC_STAR_PREFIX, - SqlTypeName.DYNAMIC_STAR) + .add("$unnest", + SqlTypeName.ANY) .nullable(true) - .buildDynamic(); + .build(); } if (type.isStruct()) { http://git-wip-us.apache.org/repos/asf/calcite/blob/68ba411e/core/src/main/java/org/apache/calcite/sql2rel/SqlToRelConverter.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/calcite/sql2rel/SqlToRelConverter.java b/core/src/main/java/org/apache/calcite/sql2rel/SqlToRelConverter.java index 9da95a5..71edcee 100644 --- a/core/src/main/java/org/apache/calcite/sql2rel/SqlToRelConverter.java +++ b/core/src/main/java/org/apache/calcite/sql2rel/SqlToRelConverter.java @@ -1956,7 +1956,7 @@ public class SqlToRelConverter { call = (SqlCall) from; convertFrom(bb, call.operand(0)); if (call.operandCount() > 2 - && bb.root instanceof Values) { + && (bb.root instanceof Values || bb.root instanceof Uncollect)) { final List<String> fieldNames = new ArrayList<>(); for (SqlNode node : Util.skip(call.getOperandList(), 2)) { fieldNames.add(((SqlIdentifier) node).getSimple()); http://git-wip-us.apache.org/repos/asf/calcite/blob/68ba411e/core/src/test/java/org/apache/calcite/test/RelOptRulesTest.java ---------------------------------------------------------------------- diff --git a/core/src/test/java/org/apache/calcite/test/RelOptRulesTest.java b/core/src/test/java/org/apache/calcite/test/RelOptRulesTest.java index 028e7eb..282449f 100644 --- a/core/src/test/java/org/apache/calcite/test/RelOptRulesTest.java +++ b/core/src/test/java/org/apache/calcite/test/RelOptRulesTest.java @@ -33,6 +33,7 @@ import org.apache.calcite.rel.RelCollationTraitDef; import org.apache.calcite.rel.RelNode; import org.apache.calcite.rel.RelRoot; import org.apache.calcite.rel.core.Aggregate; +import org.apache.calcite.rel.core.CorrelationId; import org.apache.calcite.rel.core.Intersect; import org.apache.calcite.rel.core.Join; import org.apache.calcite.rel.core.JoinRelType; @@ -40,6 +41,7 @@ import org.apache.calcite.rel.core.Minus; import org.apache.calcite.rel.core.Project; import org.apache.calcite.rel.core.RelFactories; import org.apache.calcite.rel.core.Union; +import org.apache.calcite.rel.logical.LogicalCorrelate; import org.apache.calcite.rel.logical.LogicalProject; import org.apache.calcite.rel.logical.LogicalTableModify; import org.apache.calcite.rel.logical.LogicalTableScan; @@ -76,6 +78,7 @@ import org.apache.calcite.rel.rules.JoinPushExpressionsRule; import org.apache.calcite.rel.rules.JoinPushTransitivePredicatesRule; import org.apache.calcite.rel.rules.JoinToMultiJoinRule; import org.apache.calcite.rel.rules.JoinUnionTransposeRule; +import org.apache.calcite.rel.rules.ProjectCorrelateTransposeRule; import org.apache.calcite.rel.rules.ProjectFilterTransposeRule; import org.apache.calcite.rel.rules.ProjectJoinTransposeRule; import org.apache.calcite.rel.rules.ProjectMergeRule; @@ -85,6 +88,7 @@ import org.apache.calcite.rel.rules.ProjectToCalcRule; import org.apache.calcite.rel.rules.ProjectToWindowRule; import org.apache.calcite.rel.rules.ProjectWindowTransposeRule; import org.apache.calcite.rel.rules.PruneEmptyRules; +import org.apache.calcite.rel.rules.PushProjector; import org.apache.calcite.rel.rules.ReduceExpressionsRule; import org.apache.calcite.rel.rules.SemiJoinFilterTransposeRule; import org.apache.calcite.rel.rules.SemiJoinJoinTransposeRule; @@ -103,15 +107,18 @@ import org.apache.calcite.rel.rules.UnionToDistinctRule; import org.apache.calcite.rel.rules.ValuesReduceRule; import org.apache.calcite.rel.type.RelDataType; import org.apache.calcite.rel.type.RelDataTypeFactory; +import org.apache.calcite.rex.RexCall; import org.apache.calcite.rex.RexNode; import org.apache.calcite.runtime.Hook; import org.apache.calcite.runtime.PredicateImpl; +import org.apache.calcite.sql.SemiJoinType; import org.apache.calcite.sql.SqlNode; import org.apache.calcite.sql.fun.SqlStdOperatorTable; import org.apache.calcite.sql.type.SqlTypeName; import org.apache.calcite.sql.validate.SqlValidator; import org.apache.calcite.sql2rel.SqlToRelConverter; import org.apache.calcite.tools.RelBuilder; +import org.apache.calcite.util.ImmutableBitSet; import static org.apache.calcite.plan.RelOptRule.none; import static org.apache.calcite.plan.RelOptRule.operand; @@ -174,6 +181,20 @@ import static org.junit.Assert.assertTrue; public class RelOptRulesTest extends RelOptTestBase { //~ Methods ---------------------------------------------------------------- + private final PushProjector.ExprCondition skipItem = new PushProjector.ExprCondition() { + + @Override public boolean apply(RexNode rexNode) { + return false; + } + @Override public boolean test(RexNode expr) { + if (expr instanceof RexCall) { + RexCall call = (RexCall) expr; + return "item".equalsIgnoreCase(call.getOperator().getName()); + } + return false; + } + }; + protected DiffRepository getDiffRepos() { return DiffRepository.lookup(RelOptRulesTest.class); } @@ -1051,6 +1072,138 @@ public class RelOptRulesTest extends RelOptTestBase { + "on e.ename = b.ename and e.deptno = 10"); } + @Test public void testProjectCorrelateTransposeDynamic() { + ProjectCorrelateTransposeRule customPCTrans = + new ProjectCorrelateTransposeRule(skipItem, RelFactories.LOGICAL_BUILDER); + + HepProgramBuilder programBuilder = HepProgram.builder() + .addRuleInstance(customPCTrans); + + String query = "select t1.c_nationkey, t2.a as fake_col2 " + + "from SALES.CUSTOMER as t1, " + + "unnest(t1.fake_col) as t2(a)"; + + checkPlanning( + createDynamicTester(), + null, + new HepPlanner(programBuilder.build()), + query, + true); + } + + @Test public void testProjectCorrelateTransposeRuleLeftCorrelate() { + final String sql = "SELECT e1.empno\n" + + "FROM emp e1 " + + "where exists (select empno, deptno from dept d2 where e1.deptno = d2.deptno)"; + HepProgram program = new HepProgramBuilder() + .addRuleInstance(FilterProjectTransposeRule.INSTANCE) + .addRuleInstance(ProjectFilterTransposeRule.INSTANCE) + .addRuleInstance(ProjectCorrelateTransposeRule.INSTANCE) + .build(); + sql(sql) + .withDecorrelation(false) + .expand(true) + .with(program) + .check(); + } + + @Test public void testProjectCorrelateTransposeRuleSemiCorrelate() { + RelBuilder relBuilder = RelBuilder.create(RelBuilderTest.config().build()); + RelNode left = relBuilder + .values(new String[]{"f", "f2"}, "1", "2").build(); + + CorrelationId correlationId = new CorrelationId(0); + RexNode rexCorrel = + relBuilder.getRexBuilder().makeCorrel( + left.getRowType(), + correlationId); + + RelNode right = relBuilder + .values(new String[]{"f3", "f4"}, "1", "2") + .project(relBuilder.field(0), + relBuilder.getRexBuilder() + .makeFieldAccess(rexCorrel, 0)) + .build(); + LogicalCorrelate correlate = new LogicalCorrelate(left.getCluster(), + left.getTraitSet(), left, right, correlationId, + ImmutableBitSet.of(0), SemiJoinType.SEMI); + + relBuilder.push(correlate); + RelNode relNode = relBuilder.project(relBuilder.field(0)) + .build(); + + HepProgram program = new HepProgramBuilder() + .addRuleInstance(ProjectCorrelateTransposeRule.INSTANCE) + .build(); + + HepPlanner hepPlanner = new HepPlanner(program); + hepPlanner.setRoot(relNode); + RelNode output = hepPlanner.findBestExp(); + + final String planAfter = NL + RelOptUtil.toString(output); + final DiffRepository diffRepos = getDiffRepos(); + diffRepos.assertEquals("planAfter", "${planAfter}", planAfter); + SqlToRelTestBase.assertValid(output); + } + + @Test public void testProjectCorrelateTransposeRuleAntiCorrelate() { + RelBuilder relBuilder = RelBuilder.create(RelBuilderTest.config().build()); + RelNode left = relBuilder + .values(new String[]{"f", "f2"}, "1", "2").build(); + + CorrelationId correlationId = new CorrelationId(0); + RexNode rexCorrel = + relBuilder.getRexBuilder().makeCorrel( + left.getRowType(), + correlationId); + + RelNode right = relBuilder + .values(new String[]{"f3", "f4"}, "1", "2") + .project(relBuilder.field(0), + relBuilder.getRexBuilder().makeFieldAccess(rexCorrel, 0)).build(); + LogicalCorrelate correlate = new LogicalCorrelate(left.getCluster(), + left.getTraitSet(), left, right, correlationId, + ImmutableBitSet.of(0), SemiJoinType.ANTI); + + relBuilder.push(correlate); + RelNode relNode = relBuilder.project(relBuilder.field(0)) + .build(); + + HepProgram program = new HepProgramBuilder() + .addRuleInstance(ProjectCorrelateTransposeRule.INSTANCE) + .build(); + + HepPlanner hepPlanner = new HepPlanner(program); + hepPlanner.setRoot(relNode); + RelNode output = hepPlanner.findBestExp(); + + final String planAfter = NL + RelOptUtil.toString(output); + final DiffRepository diffRepos = getDiffRepos(); + diffRepos.assertEquals("planAfter", "${planAfter}", planAfter); + SqlToRelTestBase.assertValid(output); + } + + @Test public void testProjectCorrelateTransposeWithExprCond() { + ProjectCorrelateTransposeRule customPCTrans = + new ProjectCorrelateTransposeRule(skipItem, RelFactories.LOGICAL_BUILDER); + + checkPlanning(customPCTrans, + "select t1.name, t2.ename " + + "from DEPT_NESTED as t1, " + + "unnest(t1.employees) as t2"); + } + + @Test public void testProjectCorrelateTranspose() { + ProjectCorrelateTransposeRule customPCTrans = + new ProjectCorrelateTransposeRule(PushProjector.ExprCondition.TRUE, + RelFactories.LOGICAL_BUILDER); + + checkPlanning(customPCTrans, + "select t1.name, t2.ename " + + "from DEPT_NESTED as t1, " + + "unnest(t1.employees) as t2"); + } + private static final String NOT_STRONG_EXPR = "case when e.sal < 11 then 11 else -1 * e.sal end"; http://git-wip-us.apache.org/repos/asf/calcite/blob/68ba411e/core/src/test/java/org/apache/calcite/test/RelOptTestBase.java ---------------------------------------------------------------------- diff --git a/core/src/test/java/org/apache/calcite/test/RelOptTestBase.java b/core/src/test/java/org/apache/calcite/test/RelOptTestBase.java index 70203ce..3882c84 100644 --- a/core/src/test/java/org/apache/calcite/test/RelOptTestBase.java +++ b/core/src/test/java/org/apache/calcite/test/RelOptTestBase.java @@ -60,6 +60,9 @@ abstract class RelOptTestBase extends SqlToRelTestBase { return super.createTester().withDecorrelation(false); } + protected Tester createDynamicTester() { + return getTesterWithDynamicTable(); + } /** * Checks the plan for a SQL statement before/after executing a given rule. * @@ -78,6 +81,38 @@ abstract class RelOptTestBase extends SqlToRelTestBase { } /** + * Checks the plan for a SQL statement before/after executing a given rule. + * + * @param rule Planner rule + * @param sql SQL query + */ + protected void checkPlanningDynamic( + RelOptRule rule, + String sql) { + HepProgramBuilder programBuilder = HepProgram.builder(); + programBuilder.addRuleInstance(rule); + + checkPlanning( + createDynamicTester(), + null, + new HepPlanner(programBuilder.build()), + sql); + } + /** + * Checks the plan for a SQL statement before/after executing a given rule. + * + * @param sql SQL query + */ + protected void checkPlanningDynamic( + String sql) { + checkPlanning( + createDynamicTester(), + null, + new HepPlanner(HepProgram.builder().build()), + sql); + } + + /** * Checks the plan for a SQL statement before/after executing a given * program. * http://git-wip-us.apache.org/repos/asf/calcite/blob/68ba411e/core/src/test/java/org/apache/calcite/test/SqlToRelConverterTest.java ---------------------------------------------------------------------- diff --git a/core/src/test/java/org/apache/calcite/test/SqlToRelConverterTest.java b/core/src/test/java/org/apache/calcite/test/SqlToRelConverterTest.java index bc6d6bd..a48a791 100644 --- a/core/src/test/java/org/apache/calcite/test/SqlToRelConverterTest.java +++ b/core/src/test/java/org/apache/calcite/test/SqlToRelConverterTest.java @@ -27,10 +27,8 @@ import org.apache.calcite.rel.RelRoot; import org.apache.calcite.rel.RelVisitor; import org.apache.calcite.rel.core.CorrelationId; import org.apache.calcite.rel.externalize.RelXmlWriter; -import org.apache.calcite.rel.type.RelDataType; import org.apache.calcite.rel.type.RelDataTypeFactory; import org.apache.calcite.sql.SqlExplainLevel; -import org.apache.calcite.sql.type.SqlTypeName; import org.apache.calcite.sql.validate.SqlConformance; import org.apache.calcite.sql.validate.SqlConformanceEnum; import org.apache.calcite.sql2rel.SqlToRelConverter; @@ -2497,11 +2495,25 @@ public class SqlToRelConverterTest extends SqlToRelTestBase { @Test public void testDynamicSchemaUnnest() { final String sql3 = "select t1.c_nationkey, t3.fake_col3\n" + "from SALES.CUSTOMER as t1,\n" - + "lateral (select t2.fake_col2 as fake_col3\n" + + "lateral (select t2.\"$unnest\" as fake_col3\n" + " from unnest(t1.fake_col) as t2) as t3"; sql(sql3).with(getTesterWithDynamicTable()).ok(); } + @Test public void testStarDynamicSchemaUnnest() { + final String sql3 = "select * \n" + + "from SALES.CUSTOMER as t1,\n" + + "lateral (select t2.\"$unnest\" as fake_col3\n" + + " from unnest(t1.fake_col) as t2) as t3"; + sql(sql3).with(getTesterWithDynamicTable()).ok(); + } + + @Test public void testStarDynamicSchemaUnnest2() { + final String sql3 = "select * \n" + + "from SALES.CUSTOMER as t1,\n" + + "unnest(t1.fake_col) as t2"; + sql(sql3).with(getTesterWithDynamicTable()).ok(); + } /** * Test case for Dynamic Table / Dynamic Star support * <a href="https://issues.apache.org/jira/browse/CALCITE-1150">[CALCITE-1150]</a> @@ -2607,47 +2619,6 @@ public class SqlToRelConverterTest extends SqlToRelTestBase { }); } - private Tester getTesterWithDynamicTable() { - return tester.withCatalogReaderFactory( - new Function<RelDataTypeFactory, Prepare.CatalogReader>() { - public Prepare.CatalogReader apply(RelDataTypeFactory typeFactory) { - return new MockCatalogReader(typeFactory, true) { - @Override public MockCatalogReader init() { - // CREATE SCHEMA "SALES; - // CREATE DYNAMIC TABLE "NATION" - // CREATE DYNAMIC TABLE "CUSTOMER" - - MockSchema schema = new MockSchema("SALES"); - registerSchema(schema); - - MockTable nationTable = new MockDynamicTable(this, schema.getCatalogName(), - schema.getName(), "NATION", false, 100); - registerTable(nationTable); - - MockTable customerTable = new MockDynamicTable(this, schema.getCatalogName(), - schema.getName(), "CUSTOMER", false, 100); - registerTable(customerTable); - - // CREATE TABLE "REGION" - static table with known schema. - final RelDataType intType = - typeFactory.createSqlType(SqlTypeName.INTEGER); - final RelDataType varcharType = - typeFactory.createSqlType(SqlTypeName.VARCHAR); - - MockTable regionTable = MockTable.create(this, schema, "REGION", false, 100); - regionTable.addColumn("R_REGIONKEY", intType); - regionTable.addColumn("R_NAME", varcharType); - regionTable.addColumn("R_COMMENT", varcharType); - registerTable(regionTable); - - return this; - } - // CHECKSTYLE: IGNORE 1 - }.init(); - } - }); - } - @Test public void testLarge() { SqlValidatorTest.checkLarge(400, new Function<String, Void>() { http://git-wip-us.apache.org/repos/asf/calcite/blob/68ba411e/core/src/test/java/org/apache/calcite/test/SqlToRelTestBase.java ---------------------------------------------------------------------- diff --git a/core/src/test/java/org/apache/calcite/test/SqlToRelTestBase.java b/core/src/test/java/org/apache/calcite/test/SqlToRelTestBase.java index c1b42d0..7bc3b88 100644 --- a/core/src/test/java/org/apache/calcite/test/SqlToRelTestBase.java +++ b/core/src/test/java/org/apache/calcite/test/SqlToRelTestBase.java @@ -48,6 +48,7 @@ import org.apache.calcite.sql.SqlOperatorTable; import org.apache.calcite.sql.fun.SqlStdOperatorTable; import org.apache.calcite.sql.parser.SqlParser; import org.apache.calcite.sql.type.SqlTypeFactoryImpl; +import org.apache.calcite.sql.type.SqlTypeName; import org.apache.calcite.sql.validate.SqlConformance; import org.apache.calcite.sql.validate.SqlConformanceEnum; import org.apache.calcite.sql.validate.SqlMonotonicity; @@ -103,6 +104,52 @@ public abstract class SqlToRelTestBase { SqlConformanceEnum.DEFAULT, Contexts.empty()); } + protected Tester createTester(SqlConformance conformance) { + return new TesterImpl(getDiffRepos(), false, false, true, false, + null, null, SqlToRelConverter.Config.DEFAULT, conformance, Contexts.empty()); + } + + protected Tester getTesterWithDynamicTable() { + return tester.withCatalogReaderFactory( + new Function<RelDataTypeFactory, Prepare.CatalogReader>() { + public Prepare.CatalogReader apply(RelDataTypeFactory typeFactory) { + return new MockCatalogReader(typeFactory, true) { + @Override public MockCatalogReader init() { + // CREATE SCHEMA "SALES; + // CREATE DYNAMIC TABLE "NATION" + // CREATE DYNAMIC TABLE "CUSTOMER" + + MockSchema schema = new MockSchema("SALES"); + registerSchema(schema); + + MockTable nationTable = new MockDynamicTable(this, schema.getCatalogName(), + schema.getName(), "NATION", false, 100); + registerTable(nationTable); + + MockTable customerTable = new MockDynamicTable(this, schema.getCatalogName(), + schema.getName(), "CUSTOMER", false, 100); + registerTable(customerTable); + + // CREATE TABLE "REGION" - static table with known schema. + final RelDataType intType = + typeFactory.createSqlType(SqlTypeName.INTEGER); + final RelDataType varcharType = + typeFactory.createSqlType(SqlTypeName.VARCHAR); + + MockTable regionTable = MockTable.create(this, schema, "REGION", false, 100); + regionTable.addColumn("R_REGIONKEY", intType); + regionTable.addColumn("R_NAME", varcharType); + regionTable.addColumn("R_COMMENT", varcharType); + registerTable(regionTable); + + return this; + } + // CHECKSTYLE: IGNORE 1 + }.init(); + } + }); + } + /** * Returns the default diff repository for this test, or null if there is * no repository. http://git-wip-us.apache.org/repos/asf/calcite/blob/68ba411e/core/src/test/resources/org/apache/calcite/test/RelOptRulesTest.xml ---------------------------------------------------------------------- diff --git a/core/src/test/resources/org/apache/calcite/test/RelOptRulesTest.xml b/core/src/test/resources/org/apache/calcite/test/RelOptRulesTest.xml index 36ef8ab..a5d1e56 100644 --- a/core/src/test/resources/org/apache/calcite/test/RelOptRulesTest.xml +++ b/core/src/test/resources/org/apache/calcite/test/RelOptRulesTest.xml @@ -8459,4 +8459,183 @@ LogicalProject(EMPNO=[$0], ENAME=[$1], JOB=[$2], MGR=[$3], HIREDATE=[$4], SAL=[$ ]]> </Resource> </TestCase> + <TestCase name="testProjectCorrelateTransposeWithExprCond"> + <Resource name="sql"> + <![CDATA[select t1.name, t2.ename +from DEPT_NESTED as t1, +unnest(t1.employees) as t2]]> + </Resource> + <Resource name="planBefore"> + <![CDATA[ +LogicalProject(NAME=[$1], ENAME=[$5]) + LogicalCorrelate(correlation=[$cor0], joinType=[inner], requiredColumns=[{3}]) + LogicalTableScan(table=[[CATALOG, SALES, DEPT_NESTED]]) + Uncollect + LogicalProject(EMPLOYEES=[$cor0.EMPLOYEES]) + LogicalValues(tuples=[[{ 0 }]]) +]]> + </Resource> + + <Resource name="planAfter"> + <![CDATA[ +LogicalProject(NAME=[$0], ENAME=[$2]) + LogicalCorrelate(correlation=[$cor1], joinType=[inner], requiredColumns=[{1}]) + LogicalProject(NAME=[$1], EMPLOYEES=[$3]) + LogicalTableScan(table=[[CATALOG, SALES, DEPT_NESTED]]) + LogicalProject(ENAME=[$1]) + Uncollect + LogicalProject(EMPLOYEES=[$cor1.EMPLOYEES]) + LogicalValues(tuples=[[{ 0 }]]) +]]> + </Resource> + </TestCase> + <TestCase name="testProjectCorrelateTransposeDynamic"> + <Resource name="sql"> + <![CDATA[select t1.c_nationkey, t2.fake_col2 +from SALES.CUSTOMER as t1, +unnest(t1.fake_col) as t2]]> + </Resource> + <Resource name="planBefore"> + <![CDATA[ +LogicalProject(C_NATIONKEY=[$1], FAKE_COL2=[$2]) + LogicalCorrelate(correlation=[$cor0], joinType=[inner], requiredColumns=[{0}]) + LogicalTableScan(table=[[CATALOG, SALES, CUSTOMER]]) + LogicalProject(A=[$0]) + Uncollect + LogicalProject(FAKE_COL=[$cor0.FAKE_COL]) + LogicalValues(tuples=[[{ 0 }]]) +]]> + </Resource> + <Resource name="planAfter"> + <![CDATA[ +LogicalProject(C_NATIONKEY=[$1], FAKE_COL2=[$2]) + LogicalCorrelate(correlation=[$cor0], joinType=[inner], requiredColumns=[{0}]) + LogicalTableScan(table=[[CATALOG, SALES, CUSTOMER]]) + LogicalProject(A=[$0]) + Uncollect + LogicalProject(FAKE_COL=[$cor0.FAKE_COL]) + LogicalValues(tuples=[[{ 0 }]]) +]]> + </Resource> + </TestCase> + <TestCase name="testProjectCorrelateTranspose"> + <Resource name="sql"> + <![CDATA[select t1.name, t2.ename +from DEPT_NESTED as t1, +unnest(t1.employees) as t2]]> + </Resource> + <Resource name="planBefore"> + <![CDATA[ +LogicalProject(NAME=[$1], ENAME=[$5]) + LogicalCorrelate(correlation=[$cor0], joinType=[inner], requiredColumns=[{3}]) + LogicalTableScan(table=[[CATALOG, SALES, DEPT_NESTED]]) + Uncollect + LogicalProject(EMPLOYEES=[$cor0.EMPLOYEES]) + LogicalValues(tuples=[[{ 0 }]]) +]]> + </Resource> + + <Resource name="planAfter"> + <![CDATA[ +LogicalProject(NAME=[$0], ENAME=[$2]) + LogicalCorrelate(correlation=[$cor1], joinType=[inner], requiredColumns=[{1}]) + LogicalProject(NAME=[$1], EMPLOYEES=[$3]) + LogicalTableScan(table=[[CATALOG, SALES, DEPT_NESTED]]) + LogicalProject(ENAME=[$1]) + Uncollect + LogicalProject(EMPLOYEES=[$cor1.EMPLOYEES]) + LogicalValues(tuples=[[{ 0 }]]) +]]> + </Resource> + </TestCase> + <TestCase name="testProjectCorrelateTranspose"> + <Resource name="sql"> + <![CDATA[select t1.name, t2.ename +from DEPT_NESTED as t1, +unnest(t1.employees) as t2]]> + </Resource> + <Resource name="planBefore"> + <![CDATA[ +LogicalProject(NAME=[$1], ENAME=[$5]) + LogicalCorrelate(correlation=[$cor0], joinType=[inner], requiredColumns=[{3}]) + LogicalTableScan(table=[[CATALOG, SALES, DEPT_NESTED]]) + Uncollect + LogicalProject(EMPLOYEES=[$cor0.EMPLOYEES]) + LogicalValues(tuples=[[{ 0 }]]) +]]> + </Resource> + + <Resource name="planAfter"> + <![CDATA[ +LogicalProject(NAME=[$0], ENAME=[$2]) + LogicalCorrelate(correlation=[$cor1], joinType=[inner], requiredColumns=[{1}]) + LogicalProject(NAME=[$1], EMPLOYEES=[$3]) + LogicalTableScan(table=[[CATALOG, SALES, DEPT_NESTED]]) + LogicalProject(ENAME=[$1]) + Uncollect + LogicalProject(EMPLOYEES=[$cor1.EMPLOYEES]) + LogicalValues(tuples=[[{ 0 }]]) +]]> + </Resource> + </TestCase> + <TestCase name="testProjectCorrelateTransposeRuleLeftCorrelate"> + <Resource name="planBefore"> + <![CDATA[ +LogicalProject(EMPNO=[$0]) + LogicalFilter(condition=[IS NOT NULL($9)]) + LogicalCorrelate(correlation=[$cor1], joinType=[left], requiredColumns=[{0, 7}]) + LogicalTableScan(table=[[CATALOG, SALES, EMP]]) + LogicalAggregate(group=[{}], agg#0=[MIN($0)]) + LogicalProject($f0=[true]) + LogicalProject(EMPNO=[$cor1.EMPNO], DEPTNO=[$0]) + LogicalFilter(condition=[=($cor1.DEPTNO, $0)]) + LogicalTableScan(table=[[CATALOG, SALES, DEPT]]) +]]> + </Resource> + + <Resource name="planAfter"> + <![CDATA[ +LogicalProject(EMPNO=[$0]) + LogicalFilter(condition=[IS NOT NULL($1)]) + LogicalProject(EMPNO=[$0], $f0=[$2]) + LogicalCorrelate(correlation=[$cor2], joinType=[left], requiredColumns=[{0, 1}]) + LogicalProject(EMPNO=[$0], DEPTNO=[$7]) + LogicalTableScan(table=[[CATALOG, SALES, EMP]]) + LogicalProject($f0=[$0]) + LogicalAggregate(group=[{}], agg#0=[MIN($0)]) + LogicalProject($f0=[true]) + LogicalProject(EMPNO=[$cor2.EMPNO], DEPTNO=[$0]) + LogicalFilter(condition=[=($cor2.DEPTNO, $0)]) + LogicalProject(DEPTNO=[$0]) + LogicalTableScan(table=[[CATALOG, SALES, DEPT]]) +]]> + </Resource> + </TestCase> + + <TestCase name="testProjectCorrelateTransposeRuleSemiCorrelate"> + <Resource name="planAfter"> + <![CDATA[ +LogicalCorrelate(correlation=[$cor0], joinType=[semi], requiredColumns=[{0}]) + LogicalProject(f=[$0]) + LogicalValues(tuples=[[{ '1', '2' }]]) + LogicalProject + LogicalProject(f3=[$0], $f1=[$cor0.f]) + LogicalValues(tuples=[[{ '1', '2' }]]) +]]> + </Resource> + </TestCase> + + <TestCase name="testProjectCorrelateTransposeRuleAntiCorrelate"> + <Resource name="planAfter"> + <![CDATA[ +LogicalCorrelate(correlation=[$cor0], joinType=[anti], requiredColumns=[{0}]) + LogicalProject(f=[$0]) + LogicalValues(tuples=[[{ '1', '2' }]]) + LogicalProject + LogicalProject(f3=[$0], $f1=[$cor0.f]) + LogicalValues(tuples=[[{ '1', '2' }]]) +]]> + </Resource> + </TestCase> + </Root> http://git-wip-us.apache.org/repos/asf/calcite/blob/68ba411e/core/src/test/resources/org/apache/calcite/test/SqlToRelConverterTest.xml ---------------------------------------------------------------------- diff --git a/core/src/test/resources/org/apache/calcite/test/SqlToRelConverterTest.xml b/core/src/test/resources/org/apache/calcite/test/SqlToRelConverterTest.xml index 160792a..713c2cf 100644 --- a/core/src/test/resources/org/apache/calcite/test/SqlToRelConverterTest.xml +++ b/core/src/test/resources/org/apache/calcite/test/SqlToRelConverterTest.xml @@ -750,7 +750,7 @@ lateral (select t2.fake_col2 as fake_col3 LogicalProject(C_NATIONKEY=[$1], FAKE_COL3=[$2]) LogicalCorrelate(correlation=[$cor0], joinType=[inner], requiredColumns=[{0}]) LogicalTableScan(table=[[CATALOG, SALES, CUSTOMER]]) - LogicalProject(FAKE_COL3=[ITEM($0, 'FAKE_COL2')]) + LogicalProject(FAKE_COL3=[$0]) Uncollect LogicalProject(FAKE_COL=[$cor0.FAKE_COL]) LogicalValues(tuples=[[{ 0 }]]) @@ -4543,9 +4543,10 @@ LogicalProject(DEPTNO=[$0], EMPNO=[$7]) LogicalCorrelate(correlation=[$cor0], joinType=[inner], requiredColumns=[{6}]) LogicalProject(DEPTNO=[$0], NAME=[$1], TYPE=[$2.TYPE], DESC=[$2.DESC], A=[$2.OTHERS.A], B=[$2.OTHERS.B], EMPLOYEES=[$3]) LogicalTableScan(table=[[CATALOG, SALES, DEPT_NESTED]]) - Uncollect - LogicalProject(EMPLOYEES=[$cor0.EMPLOYEES_6]) - LogicalValues(tuples=[[{ 0 }]]) + LogicalProject(EMPNO=[$0], Y=[$1], Z=[$2]) + Uncollect + LogicalProject(EMPLOYEES=[$cor0.EMPLOYEES_6]) + LogicalValues(tuples=[[{ 0 }]]) ]]> </Resource> </TestCase> @@ -5244,4 +5245,37 @@ LogicalProject(A=[$0], B=[$1]) ]]> </Resource> </TestCase> + <TestCase name="testStarDynamicSchemaUnnest"> + <Resource name="sql"> + <![CDATA[select t3.fake_q1['fake_col2'] as fake2 + from (select t2.fake_col as fake_q1 from SALES.CUSTOMER as t2) as t3]]> + </Resource> + <Resource name="plan"> + <![CDATA[ +LogicalProject(**=[$1], FAKE_COL3=[$2]) + LogicalCorrelate(correlation=[$cor0], joinType=[inner], requiredColumns=[{0}]) + LogicalTableScan(table=[[CATALOG, SALES, CUSTOMER]]) + LogicalProject(FAKE_COL3=[$0]) + Uncollect + LogicalProject(FAKE_COL=[$cor0.FAKE_COL]) + LogicalValues(tuples=[[{ 0 }]]) +]]> + </Resource> + </TestCase> + <TestCase name="testStarDynamicSchemaUnnest2"> + <Resource name="sql"> + <![CDATA[select t3.fake_q1['fake_col2'] as fake2 + from (select t2.fake_col as fake_q1 from SALES.CUSTOMER as t2) as t3]]> + </Resource> + <Resource name="plan"> + <![CDATA[ +LogicalProject(**=[$1], $unnest=[$2]) + LogicalCorrelate(correlation=[$cor0], joinType=[inner], requiredColumns=[{0}]) + LogicalTableScan(table=[[CATALOG, SALES, CUSTOMER]]) + Uncollect + LogicalProject(FAKE_COL=[$cor0.FAKE_COL]) + LogicalValues(tuples=[[{ 0 }]]) +]]> + </Resource> + </TestCase> </Root>
