From 9cfe7b12cf2b35778b5a9a7ad29fb500fbcefa42 Mon Sep 17 00:00:00 2001 From: Au_Miner <358671982@qq.com> Date: Mon, 29 Jun 2026 20:43:28 +0800 Subject: [PATCH 1/2] [FLINK-34335][table] Print query hints in RexSubQuery explain output --- .../plan/utils/RelTreeWriterImpl.scala | 167 +++++++++++++----- .../plan/hints/batch/JoinHintTestBase.java | 18 +- .../hints/batch/BroadcastJoinHintTest.xml | 83 ++++++++- .../plan/hints/batch/NestLoopJoinHintTest.xml | 83 ++++++++- .../hints/batch/ShuffleHashJoinHintTest.xml | 83 ++++++++- .../hints/batch/ShuffleMergeJoinHintTest.xml | 83 ++++++++- 6 files changed, 445 insertions(+), 72 deletions(-) diff --git a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/utils/RelTreeWriterImpl.scala b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/utils/RelTreeWriterImpl.scala index 4de3be78528ce..463d05d66b0cf 100644 --- a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/utils/RelTreeWriterImpl.scala +++ b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/utils/RelTreeWriterImpl.scala @@ -23,14 +23,17 @@ import org.apache.flink.table.planner.plan.metadata.FlinkRelMetadataQuery import org.apache.flink.table.planner.plan.nodes.physical.FlinkPhysicalRel import org.apache.flink.table.planner.plan.nodes.physical.stream._ +import org.apache.calcite.plan.RelOptUtil import org.apache.calcite.rel.RelNode +import org.apache.calcite.rel.RelWriter import org.apache.calcite.rel.core.{Aggregate, Correlate, Join, TableScan} import org.apache.calcite.rel.externalize.RelWriterImpl -import org.apache.calcite.rel.hint.Hintable +import org.apache.calcite.rel.hint.{Hintable, RelHint} +import org.apache.calcite.rex.{RexNode, RexSubQuery, RexVisitorImpl} import org.apache.calcite.sql.SqlExplainLevel import org.apache.calcite.util.Pair -import java.io.PrintWriter +import java.io.{PrintWriter, StringWriter} import java.util import java.util.concurrent.atomic.AtomicInteger @@ -139,49 +142,7 @@ class RelTreeWriterImpl( case _ => // ignore } - if (withQueryHint) { - rel match { - case _: Join | _: Correlate => - val joinHints = FlinkHints.getAllJoinHints(rel.asInstanceOf[Hintable].getHints) - if (joinHints.nonEmpty) { - printValues.add(Pair.of("joinHints", RelExplainUtil.hintsToString(joinHints))) - } - val stateTtlHints = FlinkHints.getAllStateTtlHints(rel.asInstanceOf[Hintable].getHints) - if (stateTtlHints.nonEmpty) { - printValues.add(Pair.of("stateTtlHints", RelExplainUtil.hintsToString(stateTtlHints))) - } - - case _: Aggregate | _: StreamPhysicalGroupAggregateBase => - val aggHints = - rel match { - case aggregate: Aggregate => aggregate.getHints - case _ => rel.asInstanceOf[StreamPhysicalGroupAggregateBase].hints - } - val stateTtlHints = FlinkHints.getAllStateTtlHints(aggHints) - if (stateTtlHints.nonEmpty) { - printValues.add(Pair.of("stateTtlHints", RelExplainUtil.hintsToString(stateTtlHints))) - } - case _ => // ignore - } - } - - if (withQueryBlockAlias) { - rel match { - case node: Hintable => - node match { - case _: TableScan => - // We don't need to pint hints about TableScan because TableScan will always - // print hints if exist. See more in such as LogicalTableScan#explainTerms - case _ => - val queryBlockAliasHints = FlinkHints.getQueryBlockAliasHints(node.getHints) - if (queryBlockAliasHints.nonEmpty) { - printValues.add( - Pair.of("hints", RelExplainUtil.hintsToString(queryBlockAliasHints))) - } - } - case _ => // ignore - } - } + addHintItems(rel, (term, value) => printValues.add(Pair.of(term, value))) if (withDuplicateChangesTrait) { rel match { @@ -279,6 +240,122 @@ class RelTreeWriterImpl( pw.println() } + override def item(term: String, value: AnyRef): RelWriter = { + if (withQueryHint || withQueryBlockAlias) { + value match { + case rexNode: RexNode if containsSubQuery(rexNode) => + return super.item(term, toHintAwareSubQueryString(rexNode)) + case _ => + } + } + super.item(term, value) + } + + private def containsSubQuery(node: RexNode): Boolean = { + var found = false + node.accept(new RexVisitorImpl[Void](true) { + override def visitSubQuery(subQuery: RexSubQuery): Void = { + found = true + null + } + }) + found + } + + private def addHintItems(rel: RelNode, addItem: (String, AnyRef) => Unit): Unit = { + addQueryHintItems(rel, addItem) + addQueryBlockAliasHintItems(rel, addItem) + } + + private def addQueryHintItems(rel: RelNode, addItem: (String, AnyRef) => Unit): Unit = { + if (withQueryHint) { + rel match { + case hintable: Hintable if rel.isInstanceOf[Join] || rel.isInstanceOf[Correlate] => + val hints = hintable.getHints + addJoinHintItems(hints, addItem) + addStateTtlHintItems(hints, addItem) + case aggregate: Aggregate => + addStateTtlHintItems(aggregate.getHints, addItem) + case aggregate: StreamPhysicalGroupAggregateBase => + addStateTtlHintItems(aggregate.hints, addItem) + case _ => // ignore + } + } + } + + private def addQueryBlockAliasHintItems(rel: RelNode, addItem: (String, AnyRef) => Unit): Unit = { + if (withQueryBlockAlias) { + rel match { + case _: TableScan => + // We don't need to pint hints about TableScan because TableScan will always + // print hints if exist. See more in such as LogicalTableScan#explainTerms + case hintable: Hintable => + val queryBlockAliasHints = FlinkHints.getQueryBlockAliasHints(hintable.getHints) + if (queryBlockAliasHints.nonEmpty) { + addItem("hints", RelExplainUtil.hintsToString(queryBlockAliasHints)) + } + case _ => // ignore + } + } + } + + private def addJoinHintItems( + hints: util.List[RelHint], + addItem: (String, AnyRef) => Unit): Unit = { + val joinHints = FlinkHints.getAllJoinHints(hints) + if (joinHints.nonEmpty) { + addItem("joinHints", RelExplainUtil.hintsToString(joinHints)) + } + } + + private def addStateTtlHintItems( + hints: util.List[RelHint], + addItem: (String, AnyRef) => Unit): Unit = { + val stateTtlHints = FlinkHints.getAllStateTtlHints(hints) + if (stateTtlHints.nonEmpty) { + addItem("stateTtlHints", RelExplainUtil.hintsToString(stateTtlHints)) + } + } + + /** + * Renders a {@link RexNode} to string, replacing each embedded {@link RexSubQuery}'s inner rel + * with a hint-aware flat-style plan string. + */ + private def toHintAwareSubQueryString(node: RexNode): String = { + var result = node.toString + var searchStart = 0 + node.accept(new RexVisitorImpl[Void](true) { + override def visitSubQuery(subQuery: RexSubQuery): Void = { + val withoutHints = RelOptUtil.toString(subQuery.rel) + val withHints = toHintAwareRelString(subQuery.rel) + val start = result.indexOf(withoutHints, searchStart) + if (start >= 0) { + result = + result.substring(0, start) + withHints + result.substring(start + withoutHints.length) + searchStart = start + withHints.length + } + null + } + }) + result + } + + /** + * Produces a flat-style plan string for the given rel that includes hint attributes, matching the + * format produced by {@link RelOptUtil#toString}. + */ + private def toHintAwareRelString(rel: RelNode): String = { + val sw = new StringWriter + val hintWriter = new RelWriterImpl(new PrintWriter(sw), explainLevel, withIdPrefix) { + override def done(node: RelNode): RelWriter = { + addHintItems(node, (term, value) => item(term, value)) + super.done(node) + } + } + rel.explain(hintWriter) + sw.toString + } + private def applyAdvice(rel: RelNode): Unit = { FlinkStreamPlanAnalyzers.ANALYZERS.foreach { analyzer => diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/batch/JoinHintTestBase.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/batch/JoinHintTestBase.java index 24ae0e4025d08..82730b4d84ed0 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/batch/JoinHintTestBase.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/batch/JoinHintTestBase.java @@ -51,9 +51,6 @@ * A test base for join hint. * *

TODO add test to cover legacy table source. - * - *

Notice: Join hints in sub-query will not be printed in AST, because {@code RexSubQuery} use - * 'RelOptUtil.toString(rel)' to print node and doesn't print hints about {@code LogicalJoin}. */ public abstract class JoinHintTestBase extends TableTestBase { @@ -843,6 +840,21 @@ void testJoinHintWithJoinHintInSubQuery() { verifyRelPlanByCustom(String.format(sql, buildCaseSensitiveStr(getTestSingleJoinHint()))); } + @Test + void testJoinHintWithDifferentHintsInSiblingSubQueries() { + String sql = + "select * from T1 WHERE a1 IN " + + "(select /*+ %s(T2) */ a2 from T2 join T3 on T2.a2 = T3.a3) " + + "OR a1 IN " + + "(select /*+ %s(T2) */ a2 from T2 join T3 on T2.a2 = T3.a3)"; + + verifyRelPlanByCustom( + String.format( + sql, + buildCaseSensitiveStr(getTestSingleJoinHint()), + buildCaseSensitiveStr(getOtherJoinHints().get(0)))); + } + @Test void testJoinHintWithJoinHintInCorrelateAndWithFilter() { String sql = diff --git a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/BroadcastJoinHintTest.xml b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/BroadcastJoinHintTest.xml index 9bf2749be2a02..e051056064eb9 100644 --- a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/BroadcastJoinHintTest.xml +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/BroadcastJoinHintTest.xml @@ -148,6 +148,77 @@ Calc(select=[b1, CAST(a1 AS INTEGER) AS EXPR$1]) :- Exchange(distribution=[broadcast]) : +- TableSourceScan(table=[[default_catalog, default_database, T1]], fields=[a1, b1]) +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[b2], metadata=[]]], fields=[b2]) +]]> + + + + + + + + + + + (c, 0), IS NOT NULL(i)), AND(<>(c0, 0), IS NOT NULL(i0)))]) ++- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a2)], select=[a1, b1, c, i, c0, a2, i0], build=[right]) + :- NestedLoopJoin(joinType=[InnerJoin], where=[true], select=[a1, b1, c, i, c0], build=[right], singleRowJoin=[true]) + : :- Calc(select=[a1, b1, c, i]) + : : +- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a2)], select=[a1, b1, c, a2, i], build=[right]) + : : :- NestedLoopJoin(joinType=[InnerJoin], where=[true], select=[a1, b1, c], build=[right], singleRowJoin=[true]) + : : : :- Exchange(distribution=[hash[a1]]) + : : : : +- Calc(select=[a1, b1], where=[IS NOT NULL(a1)]) + : : : : +- TableSourceScan(table=[[default_catalog, default_database, T1, filter=[]]], fields=[a1, b1]) + : : : +- Exchange(distribution=[broadcast]) + : : : +- HashAggregate(isMerge=[true], select=[Final_COUNT(count1$0) AS c]) + : : : +- Exchange(distribution=[single]) + : : : +- LocalHashAggregate(select=[Partial_COUNT(*) AS count1$0]) + : : : +- Calc(select=[a2]) + : : : +- HashJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3], isBroadcast=[true], build=[left]) + : : : :- Exchange(distribution=[broadcast]) + : : : : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) + : : : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) + : : +- Calc(select=[a2, true AS i]) + : : +- HashAggregate(isMerge=[true], groupBy=[a2], select=[a2]) + : : +- Exchange(distribution=[hash[a2]]) + : : +- LocalHashAggregate(groupBy=[a2], select=[a2]) + : : +- Calc(select=[a2]) + : : +- HashJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3], isBroadcast=[true], build=[left]) + : : :- Exchange(distribution=[broadcast]) + : : : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) + : : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) + : +- Exchange(distribution=[broadcast]) + : +- HashAggregate(isMerge=[true], select=[Final_COUNT(count1$0) AS c]) + : +- Exchange(distribution=[single]) + : +- LocalHashAggregate(select=[Partial_COUNT(*) AS count1$0]) + : +- Calc(select=[a2]) + : +- HashJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3], isBroadcast=[true], build=[left]) + : :- Exchange(distribution=[broadcast]) + : : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) + : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) + +- Calc(select=[a2, true AS i]) + +- HashAggregate(isMerge=[true], groupBy=[a2], select=[a2]) + +- Exchange(distribution=[hash[a2]]) + +- LocalHashAggregate(groupBy=[a2], select=[a2]) + +- Calc(select=[a2]) + +- HashJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3], isBroadcast=[true], build=[left]) + :- Exchange(distribution=[broadcast]) + : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) + +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) ]]> @@ -319,7 +390,7 @@ LogicalProject(a1=[$0], b1=[$1]) LogicalProject(EXPR$0=[$1]) LogicalAggregate(group=[{0}], EXPR$0=[COUNT($1)]) LogicalProject(a1=[$2], a2=[$0]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[BROADCAST inheritPath:[0, 0, 0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T1]], hints=[[[ALIAS inheritPath:[] options:[T1]]]]) })]) @@ -354,7 +425,7 @@ LogicalProject(a1=[$0], b1=[$1]) +- LogicalFilter(condition=[IN($0, { LogicalProject(a2=[$0]) LogicalFilter(condition=[=($cor0.a1, $0)]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[BROADCAST inheritPath:[0, 0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T3]], hints=[[[ALIAS inheritPath:[] options:[T3]]]]) })], variablesSet=[[$cor0]]) @@ -384,7 +455,7 @@ HashJoin(joinType=[LeftSemiJoin], where=[=(a1, a2)], select=[a1, b1], build=[rig LogicalProject(a1=[$0], b1=[$1]) +- LogicalFilter(condition=[IN($0, { LogicalProject(EXPR$0=[+($0, $cor0.a1)]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[BROADCAST inheritPath:[0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T3]], hints=[[[ALIAS inheritPath:[] options:[T3]]]]) })], variablesSet=[[$cor0]]) @@ -429,7 +500,7 @@ LogicalProject(a1=[$0], b1=[$1]) LogicalProject(a2=[$0]) LogicalSort(sort0=[$1], dir0=[ASC-nulls-first], fetch=[10]) LogicalProject(a2=[$0], a1=[$2]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[BROADCAST inheritPath:[0, 0, 0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T1]], hints=[[[ALIAS inheritPath:[] options:[T1]]]]) })]) @@ -461,7 +532,7 @@ HashJoin(joinType=[LeftSemiJoin], where=[=(a1, a2)], select=[a1, b1], isBroadcas LogicalProject(a1=[$0], b1=[$1]) +- LogicalFilter(condition=[IN($0, { LogicalProject(EXPR$0=[+($0, $cor0.a1)]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[BROADCAST inheritPath:[0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalProject(a3=[$2], b3=[$3]) LogicalJoin(condition=[=($0, $2)], joinType=[inner]) @@ -512,7 +583,7 @@ Calc(select=[a1, b1]) LogicalProject(a1=[$0], b1=[$1]) +- LogicalFilter(condition=[IN($0, { LogicalProject(a2=[$0]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[BROADCAST inheritPath:[0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T3]], hints=[[[ALIAS inheritPath:[] options:[T3]]]]) })]) diff --git a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/NestLoopJoinHintTest.xml b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/NestLoopJoinHintTest.xml index ee97a0df7bf9d..4368ff1ea713a 100644 --- a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/NestLoopJoinHintTest.xml +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/NestLoopJoinHintTest.xml @@ -148,6 +148,77 @@ Calc(select=[b1, CAST(a1 AS INTEGER) AS EXPR$1]) :- Exchange(distribution=[broadcast]) : +- TableSourceScan(table=[[default_catalog, default_database, T1]], fields=[a1, b1]) +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[b2], metadata=[]]], fields=[b2]) +]]> + + + + + + + + + + + (c, 0), IS NOT NULL(i)), AND(<>(c0, 0), IS NOT NULL(i0)))]) ++- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a2)], select=[a1, b1, c, i, c0, a2, i0], build=[right]) + :- NestedLoopJoin(joinType=[InnerJoin], where=[true], select=[a1, b1, c, i, c0], build=[right], singleRowJoin=[true]) + : :- Calc(select=[a1, b1, c, i]) + : : +- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a2)], select=[a1, b1, c, a2, i], build=[right]) + : : :- NestedLoopJoin(joinType=[InnerJoin], where=[true], select=[a1, b1, c], build=[right], singleRowJoin=[true]) + : : : :- Exchange(distribution=[hash[a1]]) + : : : : +- Calc(select=[a1, b1], where=[IS NOT NULL(a1)]) + : : : : +- TableSourceScan(table=[[default_catalog, default_database, T1, filter=[]]], fields=[a1, b1]) + : : : +- Exchange(distribution=[broadcast]) + : : : +- HashAggregate(isMerge=[true], select=[Final_COUNT(count1$0) AS c]) + : : : +- Exchange(distribution=[single]) + : : : +- LocalHashAggregate(select=[Partial_COUNT(*) AS count1$0]) + : : : +- Calc(select=[a2]) + : : : +- NestedLoopJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3], build=[left]) + : : : :- Exchange(distribution=[broadcast]) + : : : : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) + : : : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) + : : +- Calc(select=[a2, true AS i]) + : : +- HashAggregate(isMerge=[true], groupBy=[a2], select=[a2]) + : : +- Exchange(distribution=[hash[a2]]) + : : +- LocalHashAggregate(groupBy=[a2], select=[a2]) + : : +- Calc(select=[a2]) + : : +- NestedLoopJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3], build=[left]) + : : :- Exchange(distribution=[broadcast]) + : : : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) + : : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) + : +- Exchange(distribution=[broadcast]) + : +- HashAggregate(isMerge=[true], select=[Final_COUNT(count1$0) AS c]) + : +- Exchange(distribution=[single]) + : +- LocalHashAggregate(select=[Partial_COUNT(*) AS count1$0]) + : +- Calc(select=[a2]) + : +- NestedLoopJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3], build=[left]) + : :- Exchange(distribution=[broadcast]) + : : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) + : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) + +- Calc(select=[a2, true AS i]) + +- HashAggregate(isMerge=[true], groupBy=[a2], select=[a2]) + +- Exchange(distribution=[hash[a2]]) + +- LocalHashAggregate(groupBy=[a2], select=[a2]) + +- Calc(select=[a2]) + +- NestedLoopJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3], build=[left]) + :- Exchange(distribution=[broadcast]) + : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) + +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) ]]> @@ -319,7 +390,7 @@ LogicalProject(a1=[$0], b1=[$1]) LogicalProject(EXPR$0=[$1]) LogicalAggregate(group=[{0}], EXPR$0=[COUNT($1)]) LogicalProject(a1=[$2], a2=[$0]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[NEST_LOOP inheritPath:[0, 0, 0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T1]], hints=[[[ALIAS inheritPath:[] options:[T1]]]]) })]) @@ -354,7 +425,7 @@ LogicalProject(a1=[$0], b1=[$1]) +- LogicalFilter(condition=[IN($0, { LogicalProject(a2=[$0]) LogicalFilter(condition=[=($cor0.a1, $0)]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[NEST_LOOP inheritPath:[0, 0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T3]], hints=[[[ALIAS inheritPath:[] options:[T3]]]]) })], variablesSet=[[$cor0]]) @@ -384,7 +455,7 @@ HashJoin(joinType=[LeftSemiJoin], where=[=(a1, a2)], select=[a1, b1], build=[rig LogicalProject(a1=[$0], b1=[$1]) +- LogicalFilter(condition=[IN($0, { LogicalProject(EXPR$0=[+($0, $cor0.a1)]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[NEST_LOOP inheritPath:[0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T3]], hints=[[[ALIAS inheritPath:[] options:[T3]]]]) })], variablesSet=[[$cor0]]) @@ -429,7 +500,7 @@ LogicalProject(a1=[$0], b1=[$1]) LogicalProject(a2=[$0]) LogicalSort(sort0=[$1], dir0=[ASC-nulls-first], fetch=[10]) LogicalProject(a2=[$0], a1=[$2]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[NEST_LOOP inheritPath:[0, 0, 0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T1]], hints=[[[ALIAS inheritPath:[] options:[T1]]]]) })]) @@ -461,7 +532,7 @@ HashJoin(joinType=[LeftSemiJoin], where=[=(a1, a2)], select=[a1, b1], isBroadcas LogicalProject(a1=[$0], b1=[$1]) +- LogicalFilter(condition=[IN($0, { LogicalProject(EXPR$0=[+($0, $cor0.a1)]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[NEST_LOOP inheritPath:[0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalProject(a3=[$2], b3=[$3]) LogicalJoin(condition=[=($0, $2)], joinType=[inner]) @@ -512,7 +583,7 @@ Calc(select=[a1, b1]) LogicalProject(a1=[$0], b1=[$1]) +- LogicalFilter(condition=[IN($0, { LogicalProject(a2=[$0]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[NEST_LOOP inheritPath:[0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T3]], hints=[[[ALIAS inheritPath:[] options:[T3]]]]) })]) diff --git a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/ShuffleHashJoinHintTest.xml b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/ShuffleHashJoinHintTest.xml index 2e2911591ea2a..df49542d8dd25 100644 --- a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/ShuffleHashJoinHintTest.xml +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/ShuffleHashJoinHintTest.xml @@ -151,6 +151,77 @@ Calc(select=[b1, CAST(a1 AS INTEGER) AS EXPR$1]) : +- TableSourceScan(table=[[default_catalog, default_database, T1]], fields=[a1, b1]) +- Exchange(distribution=[hash[b2]]) +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[b2], metadata=[]]], fields=[b2]) +]]> + + + + + + + + + + + (c, 0), IS NOT NULL(i)), AND(<>(c0, 0), IS NOT NULL(i0)))]) ++- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a2)], select=[a1, b1, c, i, c0, a2, i0], build=[right]) + :- NestedLoopJoin(joinType=[InnerJoin], where=[true], select=[a1, b1, c, i, c0], build=[right], singleRowJoin=[true]) + : :- Calc(select=[a1, b1, c, i]) + : : +- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a2)], select=[a1, b1, c, a2, i], build=[right]) + : : :- NestedLoopJoin(joinType=[InnerJoin], where=[true], select=[a1, b1, c], build=[right], singleRowJoin=[true]) + : : : :- Exchange(distribution=[hash[a1]]) + : : : : +- Calc(select=[a1, b1], where=[IS NOT NULL(a1)]) + : : : : +- TableSourceScan(table=[[default_catalog, default_database, T1, filter=[]]], fields=[a1, b1]) + : : : +- Exchange(distribution=[broadcast]) + : : : +- HashAggregate(isMerge=[true], select=[Final_COUNT(count1$0) AS c]) + : : : +- Exchange(distribution=[single]) + : : : +- LocalHashAggregate(select=[Partial_COUNT(*) AS count1$0]) + : : : +- Calc(select=[a2]) + : : : +- HashJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3], build=[left]) + : : : :- Exchange(distribution=[hash[a2]]) + : : : : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) + : : : +- Exchange(distribution=[hash[a3]]) + : : : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) + : : +- Calc(select=[a2, true AS i]) + : : +- HashAggregate(isMerge=[false], groupBy=[a2], select=[a2]) + : : +- Calc(select=[a2]) + : : +- HashJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3], build=[left]) + : : :- Exchange(distribution=[hash[a2]]) + : : : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) + : : +- Exchange(distribution=[hash[a3]]) + : : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) + : +- Exchange(distribution=[broadcast]) + : +- HashAggregate(isMerge=[true], select=[Final_COUNT(count1$0) AS c]) + : +- Exchange(distribution=[single]) + : +- LocalHashAggregate(select=[Partial_COUNT(*) AS count1$0]) + : +- Calc(select=[a2]) + : +- HashJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3], build=[left]) + : :- Exchange(distribution=[hash[a2]]) + : : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) + : +- Exchange(distribution=[hash[a3]]) + : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) + +- Calc(select=[a2, true AS i]) + +- HashAggregate(isMerge=[false], groupBy=[a2], select=[a2]) + +- Calc(select=[a2]) + +- HashJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3], build=[left]) + :- Exchange(distribution=[hash[a2]]) + : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) + +- Exchange(distribution=[hash[a3]]) + +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) ]]> @@ -326,7 +397,7 @@ LogicalProject(a1=[$0], b1=[$1]) LogicalProject(EXPR$0=[$1]) LogicalAggregate(group=[{0}], EXPR$0=[COUNT($1)]) LogicalProject(a1=[$2], a2=[$0]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[SHUFFLE_HASH inheritPath:[0, 0, 0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T1]], hints=[[[ALIAS inheritPath:[] options:[T1]]]]) })]) @@ -360,7 +431,7 @@ LogicalProject(a1=[$0], b1=[$1]) +- LogicalFilter(condition=[IN($0, { LogicalProject(a2=[$0]) LogicalFilter(condition=[=($cor0.a1, $0)]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[SHUFFLE_HASH inheritPath:[0, 0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T3]], hints=[[[ALIAS inheritPath:[] options:[T3]]]]) })], variablesSet=[[$cor0]]) @@ -390,7 +461,7 @@ HashJoin(joinType=[LeftSemiJoin], where=[=(a1, a2)], select=[a1, b1], build=[rig LogicalProject(a1=[$0], b1=[$1]) +- LogicalFilter(condition=[IN($0, { LogicalProject(EXPR$0=[+($0, $cor0.a1)]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[SHUFFLE_HASH inheritPath:[0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T3]], hints=[[[ALIAS inheritPath:[] options:[T3]]]]) })], variablesSet=[[$cor0]]) @@ -436,7 +507,7 @@ LogicalProject(a1=[$0], b1=[$1]) LogicalProject(a2=[$0]) LogicalSort(sort0=[$1], dir0=[ASC-nulls-first], fetch=[10]) LogicalProject(a2=[$0], a1=[$2]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[SHUFFLE_HASH inheritPath:[0, 0, 0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T1]], hints=[[[ALIAS inheritPath:[] options:[T1]]]]) })]) @@ -469,7 +540,7 @@ HashJoin(joinType=[LeftSemiJoin], where=[=(a1, a2)], select=[a1, b1], isBroadcas LogicalProject(a1=[$0], b1=[$1]) +- LogicalFilter(condition=[IN($0, { LogicalProject(EXPR$0=[+($0, $cor0.a1)]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[SHUFFLE_HASH inheritPath:[0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalProject(a3=[$2], b3=[$3]) LogicalJoin(condition=[=($0, $2)], joinType=[inner]) @@ -520,7 +591,7 @@ Calc(select=[a1, b1]) LogicalProject(a1=[$0], b1=[$1]) +- LogicalFilter(condition=[IN($0, { LogicalProject(a2=[$0]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[SHUFFLE_HASH inheritPath:[0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T3]], hints=[[[ALIAS inheritPath:[] options:[T3]]]]) })]) diff --git a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/ShuffleMergeJoinHintTest.xml b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/ShuffleMergeJoinHintTest.xml index f89c07d64f3ea..62b2543ccbd39 100644 --- a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/ShuffleMergeJoinHintTest.xml +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/ShuffleMergeJoinHintTest.xml @@ -151,6 +151,77 @@ Calc(select=[b1, CAST(a1 AS INTEGER) AS EXPR$1]) : +- TableSourceScan(table=[[default_catalog, default_database, T1]], fields=[a1, b1]) +- Exchange(distribution=[hash[b2]]) +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[b2], metadata=[]]], fields=[b2]) +]]> + + + + + + + + + + + (c, 0), IS NOT NULL(i)), AND(<>(c0, 0), IS NOT NULL(i0)))]) ++- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a2)], select=[a1, b1, c, i, c0, a2, i0], build=[right]) + :- NestedLoopJoin(joinType=[InnerJoin], where=[true], select=[a1, b1, c, i, c0], build=[right], singleRowJoin=[true]) + : :- Calc(select=[a1, b1, c, i]) + : : +- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a2)], select=[a1, b1, c, a2, i], build=[right]) + : : :- NestedLoopJoin(joinType=[InnerJoin], where=[true], select=[a1, b1, c], build=[right], singleRowJoin=[true]) + : : : :- Exchange(distribution=[hash[a1]]) + : : : : +- Calc(select=[a1, b1], where=[IS NOT NULL(a1)]) + : : : : +- TableSourceScan(table=[[default_catalog, default_database, T1, filter=[]]], fields=[a1, b1]) + : : : +- Exchange(distribution=[broadcast]) + : : : +- HashAggregate(isMerge=[true], select=[Final_COUNT(count1$0) AS c]) + : : : +- Exchange(distribution=[single]) + : : : +- LocalHashAggregate(select=[Partial_COUNT(*) AS count1$0]) + : : : +- Calc(select=[a2]) + : : : +- SortMergeJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3]) + : : : :- Exchange(distribution=[hash[a2]]) + : : : : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) + : : : +- Exchange(distribution=[hash[a3]]) + : : : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) + : : +- Calc(select=[a2, true AS i]) + : : +- SortAggregate(isMerge=[false], groupBy=[a2], select=[a2]) + : : +- Calc(select=[a2]) + : : +- SortMergeJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3]) + : : :- Exchange(distribution=[hash[a2]]) + : : : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) + : : +- Exchange(distribution=[hash[a3]]) + : : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) + : +- Exchange(distribution=[broadcast]) + : +- HashAggregate(isMerge=[true], select=[Final_COUNT(count1$0) AS c]) + : +- Exchange(distribution=[single]) + : +- LocalHashAggregate(select=[Partial_COUNT(*) AS count1$0]) + : +- Calc(select=[a2]) + : +- SortMergeJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3]) + : :- Exchange(distribution=[hash[a2]]) + : : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) + : +- Exchange(distribution=[hash[a3]]) + : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) + +- Calc(select=[a2, true AS i]) + +- SortAggregate(isMerge=[false], groupBy=[a2], select=[a2]) + +- Calc(select=[a2]) + +- SortMergeJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3]) + :- Exchange(distribution=[hash[a2]]) + : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) + +- Exchange(distribution=[hash[a3]]) + +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) ]]> @@ -326,7 +397,7 @@ LogicalProject(a1=[$0], b1=[$1]) LogicalProject(EXPR$0=[$1]) LogicalAggregate(group=[{0}], EXPR$0=[COUNT($1)]) LogicalProject(a1=[$2], a2=[$0]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[SHUFFLE_MERGE inheritPath:[0, 0, 0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T1]], hints=[[[ALIAS inheritPath:[] options:[T1]]]]) })]) @@ -360,7 +431,7 @@ LogicalProject(a1=[$0], b1=[$1]) +- LogicalFilter(condition=[IN($0, { LogicalProject(a2=[$0]) LogicalFilter(condition=[=($cor0.a1, $0)]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[SHUFFLE_MERGE inheritPath:[0, 0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T3]], hints=[[[ALIAS inheritPath:[] options:[T3]]]]) })], variablesSet=[[$cor0]]) @@ -390,7 +461,7 @@ HashJoin(joinType=[LeftSemiJoin], where=[=(a1, a2)], select=[a1, b1], build=[rig LogicalProject(a1=[$0], b1=[$1]) +- LogicalFilter(condition=[IN($0, { LogicalProject(EXPR$0=[+($0, $cor0.a1)]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[SHUFFLE_MERGE inheritPath:[0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T3]], hints=[[[ALIAS inheritPath:[] options:[T3]]]]) })], variablesSet=[[$cor0]]) @@ -436,7 +507,7 @@ LogicalProject(a1=[$0], b1=[$1]) LogicalProject(a2=[$0]) LogicalSort(sort0=[$1], dir0=[ASC-nulls-first], fetch=[10]) LogicalProject(a2=[$0], a1=[$2]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[SHUFFLE_MERGE inheritPath:[0, 0, 0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T1]], hints=[[[ALIAS inheritPath:[] options:[T1]]]]) })]) @@ -469,7 +540,7 @@ HashJoin(joinType=[LeftSemiJoin], where=[=(a1, a2)], select=[a1, b1], isBroadcas LogicalProject(a1=[$0], b1=[$1]) +- LogicalFilter(condition=[IN($0, { LogicalProject(EXPR$0=[+($0, $cor0.a1)]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[SHUFFLE_MERGE inheritPath:[0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalProject(a3=[$2], b3=[$3]) LogicalJoin(condition=[=($0, $2)], joinType=[inner]) @@ -520,7 +591,7 @@ Calc(select=[a1, b1]) LogicalProject(a1=[$0], b1=[$1]) +- LogicalFilter(condition=[IN($0, { LogicalProject(a2=[$0]) - LogicalJoin(condition=[=($0, $2)], joinType=[inner]) + LogicalJoin(condition=[=($0, $2)], joinType=[inner], joinHints=[[[SHUFFLE_MERGE inheritPath:[0] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]) LogicalTableScan(table=[[default_catalog, default_database, T3]], hints=[[[ALIAS inheritPath:[] options:[T3]]]]) })]) From 31cd29dde33d5fa00096d794cc035d53595ca6de Mon Sep 17 00:00:00 2001 From: Au_Miner <358671982@qq.com> Date: Fri, 10 Jul 2026 15:05:22 +0800 Subject: [PATCH 2/2] solve comments --- .../plan/utils/RelTreeWriterImpl.scala | 48 +++++------ .../plan/hints/batch/JoinHintTestBase.java | 17 +++- .../hints/batch/BroadcastJoinHintTest.xml | 83 ++++++++++++++----- .../plan/hints/batch/NestLoopJoinHintTest.xml | 79 ++++++++++++++---- .../hints/batch/ShuffleHashJoinHintTest.xml | 83 ++++++++++++++----- .../hints/batch/ShuffleMergeJoinHintTest.xml | 83 ++++++++++++++----- .../ClearQueryBlockAliasResolverTest.xml | 45 ++++++++++ .../plan/optimize/QueryHintsResolverTest.xml | 45 ++++++++++ 8 files changed, 380 insertions(+), 103 deletions(-) diff --git a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/utils/RelTreeWriterImpl.scala b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/utils/RelTreeWriterImpl.scala index 463d05d66b0cf..15fa62e3e155c 100644 --- a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/utils/RelTreeWriterImpl.scala +++ b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/utils/RelTreeWriterImpl.scala @@ -243,25 +243,17 @@ class RelTreeWriterImpl( override def item(term: String, value: AnyRef): RelWriter = { if (withQueryHint || withQueryBlockAlias) { value match { - case rexNode: RexNode if containsSubQuery(rexNode) => - return super.item(term, toHintAwareSubQueryString(rexNode)) + case rexNode: RexNode => + toHintAwareSubQueryString(rexNode) match { + case Some(rendered) => return super.item(term, rendered) + case None => + } case _ => } } super.item(term, value) } - private def containsSubQuery(node: RexNode): Boolean = { - var found = false - node.accept(new RexVisitorImpl[Void](true) { - override def visitSubQuery(subQuery: RexSubQuery): Void = { - found = true - null - } - }) - found - } - private def addHintItems(rel: RelNode, addItem: (String, AnyRef) => Unit): Unit = { addQueryHintItems(rel, addItem) addQueryBlockAliasHintItems(rel, addItem) @@ -317,15 +309,14 @@ class RelTreeWriterImpl( } } - /** - * Renders a {@link RexNode} to string, replacing each embedded {@link RexSubQuery}'s inner rel - * with a hint-aware flat-style plan string. - */ - private def toHintAwareSubQueryString(node: RexNode): String = { + /** Renders embedded {@link RexSubQuery}s with hint-aware inner relational plans. */ + private def toHintAwareSubQueryString(node: RexNode): Option[String] = { + var found = false var result = node.toString var searchStart = 0 node.accept(new RexVisitorImpl[Void](true) { override def visitSubQuery(subQuery: RexSubQuery): Void = { + found = true val withoutHints = RelOptUtil.toString(subQuery.rel) val withHints = toHintAwareRelString(subQuery.rel) val start = result.indexOf(withoutHints, searchStart) @@ -337,18 +328,27 @@ class RelTreeWriterImpl( null } }) - result + if (found) Some(result) else None } - /** - * Produces a flat-style plan string for the given rel that includes hint attributes, matching the - * format produced by {@link RelOptUtil#toString}. - */ + /** Produces a flat-style plan string for the given rel that includes hint attributes. */ private def toHintAwareRelString(rel: RelNode): String = { val sw = new StringWriter val hintWriter = new RelWriterImpl(new PrintWriter(sw), explainLevel, withIdPrefix) { + override def item(term: String, value: AnyRef): RelWriter = { + value match { + case rexNode: RexNode => + toHintAwareSubQueryString(rexNode) match { + case Some(rendered) => super.item(term, rendered) + case None => super.item(term, value) + } + case _ => super.item(term, value) + } + } + override def done(node: RelNode): RelWriter = { - addHintItems(node, (term, value) => item(term, value)) + // RelWriterImpl stores pending attributes privately, so append hints through its item API. + addHintItems(node, (term, value) => super.item(term, value)) super.done(node) } } diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/batch/JoinHintTestBase.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/batch/JoinHintTestBase.java index 82730b4d84ed0..034b83dbddaad 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/batch/JoinHintTestBase.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/batch/JoinHintTestBase.java @@ -840,13 +840,28 @@ void testJoinHintWithJoinHintInSubQuery() { verifyRelPlanByCustom(String.format(sql, buildCaseSensitiveStr(getTestSingleJoinHint()))); } + @Test + void testJoinHintsInNestedRexSubQuery() { + String sql = + "select * from T1 where a1 in (" + + "select /*+ %s(T2) */ a2 from T2 join T3 on T2.a2 = T3.a3 " + + "where a2 in (" + + "select /*+ %s(T1) */ a1 from T1 join T3 T4 on T1.a1 = T4.a3))"; + + verifyRelPlanByCustom( + String.format( + sql, + buildCaseSensitiveStr(getTestSingleJoinHint()), + buildCaseSensitiveStr(getOtherJoinHints().get(0)))); + } + @Test void testJoinHintWithDifferentHintsInSiblingSubQueries() { String sql = "select * from T1 WHERE a1 IN " + "(select /*+ %s(T2) */ a2 from T2 join T3 on T2.a2 = T3.a3) " + "OR a1 IN " - + "(select /*+ %s(T2) */ a2 from T2 join T3 on T2.a2 = T3.a3)"; + + "(select /*+ %s(T1) */ a1 from T1 join T3 T4 on T1.a1 = T4.a3)"; verifyRelPlanByCustom( String.format( diff --git a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/BroadcastJoinHintTest.xml b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/BroadcastJoinHintTest.xml index e051056064eb9..33bdec1dac080 100644 --- a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/BroadcastJoinHintTest.xml +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/BroadcastJoinHintTest.xml @@ -153,7 +153,7 @@ Calc(select=[b1, CAST(a1 AS INTEGER) AS EXPR$1]) - + @@ -175,7 +175,7 @@ LogicalProject(a2=[$0]) (c, 0), IS NOT NULL(i)), AND(<>(c0, 0), IS NOT NULL(i0)))]) -+- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a2)], select=[a1, b1, c, i, c0, a2, i0], build=[right]) ++- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a10)], select=[a1, b1, c, i, c0, a10, i0], build=[right]) :- NestedLoopJoin(joinType=[InnerJoin], where=[true], select=[a1, b1, c, i, c0], build=[right], singleRowJoin=[true]) : :- Calc(select=[a1, b1, c, i]) : : +- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a2)], select=[a1, b1, c, a2, i], build=[right]) @@ -205,20 +205,20 @@ Calc(select=[a1, b1], where=[OR(AND(<>(c, 0), IS NOT NULL(i)), AND(<>(c0, 0), IS : +- HashAggregate(isMerge=[true], select=[Final_COUNT(count1$0) AS c]) : +- Exchange(distribution=[single]) : +- LocalHashAggregate(select=[Partial_COUNT(*) AS count1$0]) - : +- Calc(select=[a2]) - : +- HashJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3], isBroadcast=[true], build=[left]) - : :- Exchange(distribution=[broadcast]) - : : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) - : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) - +- Calc(select=[a2, true AS i]) - +- HashAggregate(isMerge=[true], groupBy=[a2], select=[a2]) - +- Exchange(distribution=[hash[a2]]) - +- LocalHashAggregate(groupBy=[a2], select=[a2]) - +- Calc(select=[a2]) - +- HashJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3], isBroadcast=[true], build=[left]) - :- Exchange(distribution=[broadcast]) - : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) - +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) + : +- Calc(select=[a1]) + : +- HashJoin(joinType=[InnerJoin], where=[=(a1, a3)], select=[a1, a3], build=[left]) + : :- Exchange(distribution=[hash[a1]]) + : : +- TableSourceScan(table=[[default_catalog, default_database, T1, project=[a1], metadata=[]]], fields=[a1], hints=[[[ALIAS options:[T1]]]]) + : +- Exchange(distribution=[hash[a3]]) + : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T4]]]]) + +- Calc(select=[a1, true AS i]) + +- HashAggregate(isMerge=[false], groupBy=[a1], select=[a1]) + +- Calc(select=[a1]) + +- HashJoin(joinType=[InnerJoin], where=[=(a1, a3)], select=[a1, a3], build=[left]) + :- Exchange(distribution=[hash[a1]]) + : +- TableSourceScan(table=[[default_catalog, default_database, T1, project=[a1], metadata=[]]], fields=[a1], hints=[[[ALIAS options:[T1]]]]) + +- Exchange(distribution=[hash[a3]]) + +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T4]]]]) ]]> @@ -1424,6 +1424,49 @@ NestedLoopJoin(joinType=[InnerJoin], where=[true], select=[a1, b1, a2, b2], buil :- Exchange(distribution=[broadcast]) : +- TableSourceScan(table=[[default_catalog, default_database, T1]], fields=[a1, b1]) +- TableSourceScan(table=[[default_catalog, default_database, T2]], fields=[a2, b2]) +]]> + + + + + + + + + + + diff --git a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/NestLoopJoinHintTest.xml b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/NestLoopJoinHintTest.xml index 4368ff1ea713a..ab4be45922469 100644 --- a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/NestLoopJoinHintTest.xml +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/NestLoopJoinHintTest.xml @@ -153,7 +153,7 @@ Calc(select=[b1, CAST(a1 AS INTEGER) AS EXPR$1]) - + @@ -175,7 +175,7 @@ LogicalProject(a2=[$0]) (c, 0), IS NOT NULL(i)), AND(<>(c0, 0), IS NOT NULL(i0)))]) -+- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a2)], select=[a1, b1, c, i, c0, a2, i0], build=[right]) ++- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a10)], select=[a1, b1, c, i, c0, a10, i0], build=[right]) :- NestedLoopJoin(joinType=[InnerJoin], where=[true], select=[a1, b1, c, i, c0], build=[right], singleRowJoin=[true]) : :- Calc(select=[a1, b1, c, i]) : : +- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a2)], select=[a1, b1, c, a2, i], build=[right]) @@ -205,20 +205,20 @@ Calc(select=[a1, b1], where=[OR(AND(<>(c, 0), IS NOT NULL(i)), AND(<>(c0, 0), IS : +- HashAggregate(isMerge=[true], select=[Final_COUNT(count1$0) AS c]) : +- Exchange(distribution=[single]) : +- LocalHashAggregate(select=[Partial_COUNT(*) AS count1$0]) - : +- Calc(select=[a2]) - : +- NestedLoopJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3], build=[left]) + : +- Calc(select=[a1]) + : +- HashJoin(joinType=[InnerJoin], where=[=(a1, a3)], select=[a1, a3], isBroadcast=[true], build=[left]) : :- Exchange(distribution=[broadcast]) - : : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) - : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) - +- Calc(select=[a2, true AS i]) - +- HashAggregate(isMerge=[true], groupBy=[a2], select=[a2]) - +- Exchange(distribution=[hash[a2]]) - +- LocalHashAggregate(groupBy=[a2], select=[a2]) - +- Calc(select=[a2]) - +- NestedLoopJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3], build=[left]) + : : +- TableSourceScan(table=[[default_catalog, default_database, T1, project=[a1], metadata=[]]], fields=[a1], hints=[[[ALIAS options:[T1]]]]) + : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T4]]]]) + +- Calc(select=[a1, true AS i]) + +- HashAggregate(isMerge=[true], groupBy=[a1], select=[a1]) + +- Exchange(distribution=[hash[a1]]) + +- LocalHashAggregate(groupBy=[a1], select=[a1]) + +- Calc(select=[a1]) + +- HashJoin(joinType=[InnerJoin], where=[=(a1, a3)], select=[a1, a3], isBroadcast=[true], build=[left]) :- Exchange(distribution=[broadcast]) - : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) - +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) + : +- TableSourceScan(table=[[default_catalog, default_database, T1, project=[a1], metadata=[]]], fields=[a1], hints=[[[ALIAS options:[T1]]]]) + +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T4]]]]) ]]> @@ -1421,6 +1421,49 @@ NestedLoopJoin(joinType=[InnerJoin], where=[true], select=[a1, b1, a2, b2], buil :- Exchange(distribution=[broadcast]) : +- TableSourceScan(table=[[default_catalog, default_database, T1]], fields=[a1, b1]) +- TableSourceScan(table=[[default_catalog, default_database, T2]], fields=[a2, b2]) +]]> + + + + + + + + + + + diff --git a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/ShuffleHashJoinHintTest.xml b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/ShuffleHashJoinHintTest.xml index df49542d8dd25..62d995b0d5300 100644 --- a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/ShuffleHashJoinHintTest.xml +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/ShuffleHashJoinHintTest.xml @@ -156,7 +156,7 @@ Calc(select=[b1, CAST(a1 AS INTEGER) AS EXPR$1]) - + @@ -178,7 +178,7 @@ LogicalProject(a2=[$0]) (c, 0), IS NOT NULL(i)), AND(<>(c0, 0), IS NOT NULL(i0)))]) -+- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a2)], select=[a1, b1, c, i, c0, a2, i0], build=[right]) ++- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a10)], select=[a1, b1, c, i, c0, a10, i0], build=[right]) :- NestedLoopJoin(joinType=[InnerJoin], where=[true], select=[a1, b1, c, i, c0], build=[right], singleRowJoin=[true]) : :- Calc(select=[a1, b1, c, i]) : : +- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a2)], select=[a1, b1, c, a2, i], build=[right]) @@ -208,20 +208,20 @@ Calc(select=[a1, b1], where=[OR(AND(<>(c, 0), IS NOT NULL(i)), AND(<>(c0, 0), IS : +- HashAggregate(isMerge=[true], select=[Final_COUNT(count1$0) AS c]) : +- Exchange(distribution=[single]) : +- LocalHashAggregate(select=[Partial_COUNT(*) AS count1$0]) - : +- Calc(select=[a2]) - : +- HashJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3], build=[left]) - : :- Exchange(distribution=[hash[a2]]) - : : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) - : +- Exchange(distribution=[hash[a3]]) - : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) - +- Calc(select=[a2, true AS i]) - +- HashAggregate(isMerge=[false], groupBy=[a2], select=[a2]) - +- Calc(select=[a2]) - +- HashJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3], build=[left]) - :- Exchange(distribution=[hash[a2]]) - : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) - +- Exchange(distribution=[hash[a3]]) - +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) + : +- Calc(select=[a1]) + : +- HashJoin(joinType=[InnerJoin], where=[=(a1, a3)], select=[a1, a3], isBroadcast=[true], build=[left]) + : :- Exchange(distribution=[broadcast]) + : : +- TableSourceScan(table=[[default_catalog, default_database, T1, project=[a1], metadata=[]]], fields=[a1], hints=[[[ALIAS options:[T1]]]]) + : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T4]]]]) + +- Calc(select=[a1, true AS i]) + +- HashAggregate(isMerge=[true], groupBy=[a1], select=[a1]) + +- Exchange(distribution=[hash[a1]]) + +- LocalHashAggregate(groupBy=[a1], select=[a1]) + +- Calc(select=[a1]) + +- HashJoin(joinType=[InnerJoin], where=[=(a1, a3)], select=[a1, a3], isBroadcast=[true], build=[left]) + :- Exchange(distribution=[broadcast]) + : +- TableSourceScan(table=[[default_catalog, default_database, T1, project=[a1], metadata=[]]], fields=[a1], hints=[[[ALIAS options:[T1]]]]) + +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T4]]]]) ]]> @@ -1446,6 +1446,49 @@ NestedLoopJoin(joinType=[InnerJoin], where=[true], select=[a1, b1, a2, b2], buil :- Exchange(distribution=[broadcast]) : +- TableSourceScan(table=[[default_catalog, default_database, T1]], fields=[a1, b1]) +- TableSourceScan(table=[[default_catalog, default_database, T2]], fields=[a2, b2]) +]]> + + + + + + + + + + + diff --git a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/ShuffleMergeJoinHintTest.xml b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/ShuffleMergeJoinHintTest.xml index 62b2543ccbd39..25ed8ba8e68ad 100644 --- a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/ShuffleMergeJoinHintTest.xml +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/hints/batch/ShuffleMergeJoinHintTest.xml @@ -156,7 +156,7 @@ Calc(select=[b1, CAST(a1 AS INTEGER) AS EXPR$1]) - + @@ -178,7 +178,7 @@ LogicalProject(a2=[$0]) (c, 0), IS NOT NULL(i)), AND(<>(c0, 0), IS NOT NULL(i0)))]) -+- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a2)], select=[a1, b1, c, i, c0, a2, i0], build=[right]) ++- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a10)], select=[a1, b1, c, i, c0, a10, i0], build=[right]) :- NestedLoopJoin(joinType=[InnerJoin], where=[true], select=[a1, b1, c, i, c0], build=[right], singleRowJoin=[true]) : :- Calc(select=[a1, b1, c, i]) : : +- HashJoin(joinType=[LeftOuterJoin], where=[=(a1, a2)], select=[a1, b1, c, a2, i], build=[right]) @@ -208,20 +208,20 @@ Calc(select=[a1, b1], where=[OR(AND(<>(c, 0), IS NOT NULL(i)), AND(<>(c0, 0), IS : +- HashAggregate(isMerge=[true], select=[Final_COUNT(count1$0) AS c]) : +- Exchange(distribution=[single]) : +- LocalHashAggregate(select=[Partial_COUNT(*) AS count1$0]) - : +- Calc(select=[a2]) - : +- SortMergeJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3]) - : :- Exchange(distribution=[hash[a2]]) - : : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) - : +- Exchange(distribution=[hash[a3]]) - : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) - +- Calc(select=[a2, true AS i]) - +- SortAggregate(isMerge=[false], groupBy=[a2], select=[a2]) - +- Calc(select=[a2]) - +- SortMergeJoin(joinType=[InnerJoin], where=[=(a2, a3)], select=[a2, a3]) - :- Exchange(distribution=[hash[a2]]) - : +- TableSourceScan(table=[[default_catalog, default_database, T2, project=[a2], metadata=[]]], fields=[a2], hints=[[[ALIAS options:[T2]]]]) - +- Exchange(distribution=[hash[a3]]) - +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T3]]]]) + : +- Calc(select=[a1]) + : +- HashJoin(joinType=[InnerJoin], where=[=(a1, a3)], select=[a1, a3], isBroadcast=[true], build=[left]) + : :- Exchange(distribution=[broadcast]) + : : +- TableSourceScan(table=[[default_catalog, default_database, T1, project=[a1], metadata=[]]], fields=[a1], hints=[[[ALIAS options:[T1]]]]) + : +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T4]]]]) + +- Calc(select=[a1, true AS i]) + +- HashAggregate(isMerge=[true], groupBy=[a1], select=[a1]) + +- Exchange(distribution=[hash[a1]]) + +- LocalHashAggregate(groupBy=[a1], select=[a1]) + +- Calc(select=[a1]) + +- HashJoin(joinType=[InnerJoin], where=[=(a1, a3)], select=[a1, a3], isBroadcast=[true], build=[left]) + :- Exchange(distribution=[broadcast]) + : +- TableSourceScan(table=[[default_catalog, default_database, T1, project=[a1], metadata=[]]], fields=[a1], hints=[[[ALIAS options:[T1]]]]) + +- TableSourceScan(table=[[default_catalog, default_database, T3, project=[a3], metadata=[]]], fields=[a3], hints=[[[ALIAS options:[T4]]]]) ]]> @@ -1446,6 +1446,49 @@ NestedLoopJoin(joinType=[InnerJoin], where=[true], select=[a1, b1, a2, b2], buil :- Exchange(distribution=[broadcast]) : +- TableSourceScan(table=[[default_catalog, default_database, T1]], fields=[a1, b1]) +- TableSourceScan(table=[[default_catalog, default_database, T2]], fields=[a2, b2]) +]]> + + + + + + + + + + + diff --git a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/optimize/ClearQueryBlockAliasResolverTest.xml b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/optimize/ClearQueryBlockAliasResolverTest.xml index 75a756dc770ce..a4962235c60c4 100644 --- a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/optimize/ClearQueryBlockAliasResolverTest.xml +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/optimize/ClearQueryBlockAliasResolverTest.xml @@ -91,6 +91,28 @@ LogicalProject(b1=[$1], EXPR$1=[CAST($0):INTEGER]), rowType=[RecordType(VARCHAR( +- LogicalJoin(condition=[=($1, $3)], joinType=[inner], joinHints=[[[BROADCAST options:[LEFT]]]]), rowType=[RecordType(BIGINT a1, VARCHAR(2147483647) b1, BIGINT a2, VARCHAR(2147483647) b2)] :- LogicalTableScan(table=[[default_catalog, default_database, T1]]), rowType=[RecordType(BIGINT a1, VARCHAR(2147483647) b1)] +- LogicalTableScan(table=[[default_catalog, default_database, T2]]), rowType=[RecordType(BIGINT a2, VARCHAR(2147483647) b2)] +]]> + + + + + + + + @@ -755,6 +777,29 @@ LogicalProject(a1=[$0], b1=[$1], a2=[$2], b2=[$3]), rowType=[RecordType(BIGINT a +- LogicalJoin(condition=[true], joinType=[inner], joinHints=[[[BROADCAST options:[LEFT]]]]), rowType=[RecordType(BIGINT a1, VARCHAR(2147483647) b1, BIGINT a2, VARCHAR(2147483647) b2)] :- LogicalTableScan(table=[[default_catalog, default_database, T1]]), rowType=[RecordType(BIGINT a1, VARCHAR(2147483647) b1)] +- LogicalTableScan(table=[[default_catalog, default_database, T2]]), rowType=[RecordType(BIGINT a2, VARCHAR(2147483647) b2)] +]]> + + + + + + + + diff --git a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/optimize/QueryHintsResolverTest.xml b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/optimize/QueryHintsResolverTest.xml index 24db25259693a..bf0bc26b7c249 100644 --- a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/optimize/QueryHintsResolverTest.xml +++ b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/optimize/QueryHintsResolverTest.xml @@ -91,6 +91,28 @@ LogicalProject(b1=[$1], EXPR$1=[CAST($0):INTEGER]), rowType=[RecordType(VARCHAR( +- LogicalJoin(condition=[=($1, $3)], joinType=[inner], joinHints=[[[BROADCAST options:[LEFT]]]]), rowType=[RecordType(BIGINT a1, VARCHAR(2147483647) b1, BIGINT a2, VARCHAR(2147483647) b2)] :- LogicalTableScan(table=[[default_catalog, default_database, T1]], hints=[[[ALIAS inheritPath:[] options:[T1]]]]), rowType=[RecordType(BIGINT a1, VARCHAR(2147483647) b1)] +- LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]), rowType=[RecordType(BIGINT a2, VARCHAR(2147483647) b2)] +]]> + + + + + + + + @@ -755,6 +777,29 @@ LogicalProject(a1=[$0], b1=[$1], a2=[$2], b2=[$3]), rowType=[RecordType(BIGINT a +- LogicalJoin(condition=[true], joinType=[inner], joinHints=[[[BROADCAST options:[LEFT]]]]), rowType=[RecordType(BIGINT a1, VARCHAR(2147483647) b1, BIGINT a2, VARCHAR(2147483647) b2)] :- LogicalTableScan(table=[[default_catalog, default_database, T1]], hints=[[[ALIAS inheritPath:[] options:[T1]]]]), rowType=[RecordType(BIGINT a1, VARCHAR(2147483647) b1)] +- LogicalTableScan(table=[[default_catalog, default_database, T2]], hints=[[[ALIAS inheritPath:[] options:[T2]]]]), rowType=[RecordType(BIGINT a2, VARCHAR(2147483647) b2)] +]]> + + + + + + + +