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 4de3be78528ce6..15fa62e3e155c3 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 =>
+ toHintAwareSubQueryString(rexNode) match {
+ case Some(rendered) => return super.item(term, rendered)
+ case None =>
+ }
+ case _ =>
+ }
+ }
+ super.item(term, value)
+ }
+
+ 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 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)
+ if (start >= 0) {
+ result =
+ result.substring(0, start) + withHints + result.substring(start + withoutHints.length)
+ searchStart = start + withHints.length
+ }
+ null
+ }
+ })
+ if (found) Some(result) else None
+ }
+
+ /** 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 = {
+ // RelWriterImpl stores pending attributes privately, so append hints through its item API.
+ addHintItems(node, (term, value) => super.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 24ae0e4025d08d..034b83dbddaada 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,36 @@ 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(T1) */ a1 from T1 join T3 T4 on T1.a1 = T4.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 9bf2749be2a022..33bdec1dac0809 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, 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])
+ : : :- 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=[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]]]])
]]>
@@ -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]]]])
})])
@@ -1353,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 ee97a0df7bf9d9..ab4be45922469c 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, 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])
+ : : :- 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=[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]]]])
]]>
@@ -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]]]])
})])
@@ -1350,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 2e2911591ea2a3..62d995b0d53009 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, 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])
+ : : :- 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=[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]]]])
]]>
@@ -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]]]])
})])
@@ -1375,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 f89c07d64f3eab..25ed8ba8e68ade 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, 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])
+ : : :- 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=[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]]]])
]]>
@@ -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]]]])
})])
@@ -1375,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 75a756dc770ce4..a4962235c60c40 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 24db25259693a5..bf0bc26b7c2493 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)]
+]]>
+
+
+
+
+
+
+
+