diff --git a/docs/content.zh/docs/sql/reference/dml/insert.md b/docs/content.zh/docs/sql/reference/dml/insert.md index 0d879c1c91a2fe..f8af55b8c90c91 100644 --- a/docs/content.zh/docs/sql/reference/dml/insert.md +++ b/docs/content.zh/docs/sql/reference/dml/insert.md @@ -314,6 +314,11 @@ This check is controlled by the configuration option `table.exec.sink.require-on Alternatively, if you do not need consistency guarantees for conflicting keys, you can disable the sink upsert materializer entirely by setting `table.exec.sink.upsert-materialize` to `NONE`. This removes the materializer operator from the pipeline, so no buffering, compaction, or conflict resolution is performed. Records are passed directly to the sink as they arrive. +### Materialized tables + +A materialized table has no syntax for an explicit `ON CONFLICT` clause, so it always falls back to the implicit `DO DEDUPLICATE` strategy when its primary key differs from its query's upsert key. Setting `table.exec.sink.materialized-table-forces-on-conflict-error` (default: `false`) makes it default to `DO ERROR` instead, avoiding `DO DEDUPLICATE`'s higher state cost. Since `DO ERROR` requires a watermark on every source, this only applies when all sources already have one - otherwise the table keeps the `DO DEDUPLICATE` default so its creation never fails because of this option. This option adds `ON CONFLICT DO ERROR` independently from `table.exec.sink.require-on-conflict`. Optimally, when working with Materialized Tables, disable `table.exec.sink.require-on-conflict` to avoid confusion. If you enable both options, all your source tables will have to define a watermark strategy. + + ### Syntax ```sql diff --git a/docs/content/docs/sql/reference/dml/insert.md b/docs/content/docs/sql/reference/dml/insert.md index 11bc66632feec7..0f4f33df85d3e5 100644 --- a/docs/content/docs/sql/reference/dml/insert.md +++ b/docs/content/docs/sql/reference/dml/insert.md @@ -324,6 +324,12 @@ This check is controlled by the configuration option `table.exec.sink.require-on Alternatively, if you do not need consistency guarantees for conflicting keys, you can disable the sink upsert materializer entirely by setting `table.exec.sink.upsert-materialize` to `NONE`. This removes the materializer operator from the pipeline, so no buffering, compaction, or conflict resolution is performed. Records are passed directly to the sink as they arrive. +### Materialized tables + +A materialized table has no syntax for an explicit `ON CONFLICT` clause, so it always falls back to the implicit `DO DEDUPLICATE` strategy when its primary key differs from its query's upsert key. Setting `table.exec.sink.materialized-table-forces-on-conflict-error` (default: `false`) makes it default to `DO ERROR` instead, avoiding `DO DEDUPLICATE`'s higher state cost. Since `DO ERROR` requires a watermark on every source, this only applies when all sources already have one - otherwise the table keeps the `DO DEDUPLICATE` default so its creation never fails because of this option. This option adds `ON CONFLICT DO ERROR` independently from `table.exec.sink.require-on-conflict`. Optimally, when working with Materialized Tables, disable `table.exec.sink.require-on-conflict` to avoid confusion. If you enable both options, all your source tables will have to define a watermark strategy. + + + ### Syntax ```sql diff --git a/docs/layouts/shortcodes/generated/execution_config_configuration.html b/docs/layouts/shortcodes/generated/execution_config_configuration.html index e71bf3d7673a92..63425656bedf21 100644 --- a/docs/layouts/shortcodes/generated/execution_config_configuration.html +++ b/docs/layouts/shortcodes/generated/execution_config_configuration.html @@ -260,6 +260,12 @@

Enum

In order to minimize the distributed disorder problem when writing data into table with primary keys that many users suffers. FLINK will auto add a keyed shuffle by default when the sink parallelism differs from upstream operator and sink parallelism is not 1. This works only when the upstream ensures the multi-records' order on the primary key, if not, the added shuffle can not solve the problem (In this situation, a more proper way is to consider the deduplicate operation for the source firstly or use an upsert source with primary key definition which truly reflect the records evolution).
By default, the keyed shuffle will be added when the sink's parallelism differs from upstream operator. You can set to no shuffle(NONE) or force shuffle(FORCE).

Possible values: + +
table.exec.sink.materialized-table-forces-on-conflict-error

Streaming + false + Boolean + When enabled, a materialized table whose declared primary key differs from its query's upsert key defaults to ON CONFLICT DO ERROR instead of the state-heavy deduplicating materializer.

Independent of table.exec.sink.require-on-conflict, which controls whether an explicit ON CONFLICT clause is required elsewhere. +
table.exec.sink.nested-constraint-enforcer

Batch Streaming IGNORE diff --git a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/config/ExecutionConfigOptions.java b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/config/ExecutionConfigOptions.java index a1ed3ec97ff588..893bf76461a563 100644 --- a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/config/ExecutionConfigOptions.java +++ b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/config/ExecutionConfigOptions.java @@ -277,6 +277,27 @@ public class ExecutionConfigOptions { + "results in certain streaming scenarios.") .build()); + @Documentation.TableOption(execMode = Documentation.ExecMode.STREAMING) + public static final ConfigOption + TABLE_EXEC_SINK_MATERIALIZED_TABLE_FORCES_ON_CONFLICT_ERROR = + key("table.exec.sink.materialized-table-forces-on-conflict-error") + .booleanType() + .defaultValue(false) + .withDescription( + Description.builder() + .text( + "When enabled, a materialized table whose declared primary key " + + "differs from its query's upsert key defaults to ON CONFLICT " + + "DO ERROR instead of the state-heavy deduplicating " + + "materializer.") + .linebreak() + .linebreak() + .text( + "Independent of table.exec.sink.require-on-conflict, which " + + "controls whether an explicit ON CONFLICT clause is " + + "required elsewhere.") + .build()); + // ------------------------------------------------------------------------ // Sort Options // ------------------------------------------------------------------------ diff --git a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/connectors/DynamicSinkUtils.java b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/connectors/DynamicSinkUtils.java index cb579c644ca511..4aee0131323e2a 100644 --- a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/connectors/DynamicSinkUtils.java +++ b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/connectors/DynamicSinkUtils.java @@ -83,6 +83,7 @@ import org.apache.calcite.rel.RelNode; import org.apache.calcite.rel.core.Project; import org.apache.calcite.rel.core.TableModify; +import org.apache.calcite.rel.core.TableScan; import org.apache.calcite.rel.hint.RelHint; import org.apache.calcite.rel.logical.LogicalFilter; import org.apache.calcite.rel.logical.LogicalTableModify; @@ -263,6 +264,10 @@ public static RelNode convertCreateTableAsToRel( * AS) into a {@link RelNode} that writes into the table, adding helper projections if * necessary. The caller resolves the target table to its {@link ResolvedCatalogTable} form (the * new definition for an alter, the created definition for a create). + * + *

The table defaults to {@code ON CONFLICT DO ERROR} when {@link + * ExecutionConfigOptions#TABLE_EXEC_SINK_MATERIALIZED_TABLE_FORCES_ON_CONFLICT_ERROR} is set + * and every source has a watermark; otherwise it keeps the deduplicating default. */ public static RelNode convertMaterializedTableAsToRel( FlinkRelBuilder relBuilder, @@ -272,10 +277,19 @@ public static RelNode convertMaterializedTableAsToRel( ResolvedCatalogTable resolvedTable, Map staticPartitions, boolean isOverwrite, - DynamicTableSink sink) { + DynamicTableSink sink, + ReadableConfig configuration) { final ContextResolvedTable contextResolvedTable = ContextResolvedTable.permanent(identifier, catalog, resolvedTable); + final InsertConflictStrategy conflictStrategy = + configuration.get( + ExecutionConfigOptions + .TABLE_EXEC_SINK_MATERIALIZED_TABLE_FORCES_ON_CONFLICT_ERROR) + && allSourcesHaveWatermarks(input) + ? InsertConflictStrategy.error() + : null; + return convertSinkToRel( relBuilder, input, @@ -285,7 +299,28 @@ public static RelNode convertMaterializedTableAsToRel( null, isOverwrite, sink, - null); + conflictStrategy); + } + + /** Whether every table scanned by {@code rel} declares a watermark. */ + private static boolean allSourcesHaveWatermarks(RelNode rel) { + if (rel instanceof TableScan) { + final TableSourceTable table = + ((TableScan) rel).getTable().unwrap(TableSourceTable.class); + // An unresolvable scan can't be proven to lack a watermark, so it isn't treated as a + // reason to fall back - matching the same convention used for the physical-tree check. + return table == null + || !table.contextResolvedTable() + .getResolvedSchema() + .getWatermarkSpecs() + .isEmpty(); + } + for (RelNode input : rel.getInputs()) { + if (!allSourcesHaveWatermarks(input)) { + return false; + } + } + return true; } private static RelNode convertSinkToRel( diff --git a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/delegation/PlannerBase.scala b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/delegation/PlannerBase.scala index e4f76fd1990df8..2c16b860be778d 100644 --- a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/delegation/PlannerBase.scala +++ b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/delegation/PlannerBase.scala @@ -615,7 +615,8 @@ abstract class PlannerBase( resolvedTable, staticPartitions, false, - tableSink) + tableSink, + getTableConfig) } protected def createSerdeContext: SerdeContext = { diff --git a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/optimize/program/FlinkChangelogModeInferenceProgram.scala b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/optimize/program/FlinkChangelogModeInferenceProgram.scala index 8d22466c3eec02..1823065059e02a 100644 --- a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/optimize/program/FlinkChangelogModeInferenceProgram.scala +++ b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/optimize/program/FlinkChangelogModeInferenceProgram.scala @@ -1249,9 +1249,25 @@ class FlinkChangelogModeInferenceProgram extends FlinkOptimizeProgram[StreamOpti if (requireOnConflict && upsertKeyDiffersFromPk && sink.conflictStrategy == null) { val pkNames = sink.getPrimaryKeyNames val upsertKeyNames = sink.getUpsertKeyNames + val identifier = sink.contextResolvedTable.getIdentifier.asSummaryString + val mtOnConflictDoError = tableConfig.get( + ExecutionConfigOptions.TABLE_EXEC_SINK_MATERIALIZED_TABLE_FORCES_ON_CONFLICT_ERROR) + val materializedTableHint = + if (mtOnConflictDoError) { + // When both requireOnConflict and mtOnConflictDoError are enabled, every source + // must declare a watermark for ON CONFLICT DO ERROR to apply. + val sourcesWithoutWatermarks = new java.util.ArrayList[String]() + collectSourcesWithoutWatermarks(sink.getInput, sourcesWithoutWatermarks) + " If this is a materialized table: table.exec.sink.materialized-table-forces-on-" + + "conflict-error is enabled, but a watermark is missing on: " + + s"${sourcesWithoutWatermarks.toArray.mkString(", ")}. Add one to every source to " + + "satisfy both options." + } else { + "" + } throw new ValidationException( - "The query has an upsert key that differs from the primary key of the sink table " + - s"'${sink.contextResolvedTable.getIdentifier.asSummaryString}'. " + + s"The query has an upsert key that differs from the primary key of the sink table " + + s"'$identifier'. " + s"Primary key: $pkNames, upsert key: $upsertKeyNames. " + "This can lead to non-deterministic results when multiple records with different " + "upsert keys map to the same primary key. " + @@ -1259,7 +1275,8 @@ class FlinkChangelogModeInferenceProgram extends FlinkOptimizeProgram[StreamOpti "ON CONFLICT DO DEDUPLICATE (update to the latest record, state intensive, since we" + " need to keep the entire history), or " + "ON CONFLICT DO ERROR (fail on conflict), or " + - "ON CONFLICT DO NOTHING (keep first record).") + "ON CONFLICT DO NOTHING (keep first record)." + + materializedTableHint) } } diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ExplainTest.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ExplainTest.java index e6b4024fae91f7..52ae0e17feb82e 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ExplainTest.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ExplainTest.java @@ -20,6 +20,8 @@ import org.apache.flink.configuration.Configuration; import org.apache.flink.table.api.TableConfig; +import org.apache.flink.table.api.ValidationException; +import org.apache.flink.table.api.config.ExecutionConfigOptions; import org.apache.flink.table.api.config.TableConfigOptions; import org.apache.flink.table.planner.utils.StreamTableTestUtil; import org.apache.flink.table.planner.utils.TableTestBase; @@ -34,6 +36,7 @@ import java.nio.file.Paths; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; /** Tests for EXPLAIN statements. */ public class ExplainTest extends TableTestBase { @@ -78,6 +81,149 @@ void testExplainCreateMaterializedTable() { + " MyTable"); } + @Test + void testExplainCreateMaterializedTableDefaultsToOnConflictErrorWhenEnabled() { + util.getTableEnv() + .getConfig() + .set( + ExecutionConfigOptions + .TABLE_EXEC_SINK_MATERIALIZED_TABLE_FORCES_ON_CONFLICT_ERROR, + true); + util.getTableEnv() + .executeSql( + "CREATE TABLE MyWatermarkedTable (\n" + + " a INT,\n" + + " c STRING,\n" + + " ts TIMESTAMP(3),\n" + + " WATERMARK FOR ts AS ts\n" + + ") WITH (\n" + + " 'connector' = 'values'\n" + + ")"); + // Upsert key is (a, c) from GROUP BY; declared PK is (a) alone - a conflict. + verifyExplain( + "CREATE MATERIALIZED TABLE MyMTOnConflictTable (\n" + + " a INT,\n" + + " cnt BIGINT,\n" + + " PRIMARY KEY (cnt) NOT ENFORCED\n" + + ") WITH (\n" + + " 'connector' = 'values',\n" + + " 'sink-insert-only' = 'false'\n" + + ") AS\n" + + " SELECT a, COUNT(*) AS cnt FROM MyWatermarkedTable GROUP BY a, c"); + } + + @Test + void testExplainCreateMaterializedTableNoDuplicatesOnConflictErrorWhenEnabled() { + util.getTableEnv() + .getConfig() + .set( + ExecutionConfigOptions + .TABLE_EXEC_SINK_MATERIALIZED_TABLE_FORCES_ON_CONFLICT_ERROR, + true); + // Upsert key (a) matches the declared PK (a): no materializer should be inserted, + // flag or no flag. COALESCE keeps the grouping key NOT NULL for the PK column. + verifyExplain( + "CREATE MATERIALIZED TABLE MyMTNoConflictTable (\n" + + " a INT,\n" + + " cnt BIGINT,\n" + + " PRIMARY KEY (a) NOT ENFORCED\n" + + ") WITH (\n" + + " 'connector' = 'values',\n" + + " 'sink-insert-only' = 'false'\n" + + ") AS\n" + + " SELECT COALESCE(a, 0) AS a, COUNT(*) AS cnt FROM MyTable" + + " GROUP BY COALESCE(a, 0)"); + } + + @Test + void testExplainCreateMaterializedTableKeepsOldDefaultWhenSourceHasNoWatermark() { + util.getTableEnv() + .getConfig() + .set( + ExecutionConfigOptions + .TABLE_EXEC_SINK_MATERIALIZED_TABLE_FORCES_ON_CONFLICT_ERROR, + true); + // Isolate from table.exec.sink.require-on-conflict as above. + util.getTableEnv() + .getConfig() + .set(ExecutionConfigOptions.TABLE_EXEC_SINK_REQUIRE_ON_CONFLICT, false); + util.getTableEnv() + .executeSql( + "CREATE TABLE MyWatermarkedTable (\n" + + " a INT,\n" + + " c STRING,\n" + + " ts TIMESTAMP(3),\n" + + " WATERMARK FOR ts AS ts\n" + + ") WITH (\n" + + " 'connector' = 'values'\n" + + ")"); + // A join across a watermarked and an unwatermarked source: not every source has a + // watermark, so the fallback must still apply even though one branch alone would pass. + verifyExplain( + "CREATE MATERIALIZED TABLE MyMTJoinedSourcesTable (\n" + + " a INT,\n" + + " cnt BIGINT,\n" + + " PRIMARY KEY (cnt) NOT ENFORCED\n" + + ") WITH (\n" + + " 'connector' = 'values',\n" + + " 'sink-insert-only' = 'false'\n" + + ") AS\n" + + " SELECT w.a, COUNT(*) AS cnt FROM MyWatermarkedTable w" + + " JOIN MyTable m ON w.a = m.a GROUP BY w.a, w.c"); + } + + @Test + void testExplainCreateMaterializedTableKeepsOldDefaultWhenDisabled() { + // Isolate from table.exec.sink.require-on-conflict as above. + util.getTableEnv() + .getConfig() + .set(ExecutionConfigOptions.TABLE_EXEC_SINK_REQUIRE_ON_CONFLICT, false); + verifyExplain( + "CREATE MATERIALIZED TABLE MyMTNoConflictDefaultTable (\n" + + " a INT,\n" + + " cnt BIGINT,\n" + + " PRIMARY KEY (cnt) NOT ENFORCED\n" + + ") WITH (\n" + + " 'connector' = 'values',\n" + + " 'sink-insert-only' = 'false'\n" + + ") AS\n" + + " SELECT a, COUNT(*) AS cnt FROM MyTable GROUP BY a, c"); + } + + @Test + void testExplainCreateMaterializedTableErrorMentionsWatermarkWhenBothOptionsEnabled() { + util.getTableEnv() + .getConfig() + .set( + ExecutionConfigOptions + .TABLE_EXEC_SINK_MATERIALIZED_TABLE_FORCES_ON_CONFLICT_ERROR, + true); + util.getTableEnv() + .getConfig() + .set(ExecutionConfigOptions.TABLE_EXEC_SINK_REQUIRE_ON_CONFLICT, true); + // With both options on and no watermark, the message also points at the missing + // watermark - the actionable fix for a materialized table, which has no ON CONFLICT + // syntax to satisfy the generic advice above it. + assertThatThrownBy( + () -> + util.getTableEnv() + .explainSql( + "CREATE MATERIALIZED TABLE MyMTBothOptionsTable (\n" + + " a INT,\n" + + " cnt BIGINT,\n" + + " PRIMARY KEY (cnt) NOT ENFORCED\n" + + ") WITH (\n" + + " 'connector' = 'values',\n" + + " 'sink-insert-only' = 'false'\n" + + ") AS\n" + + " SELECT a, COUNT(*) AS cnt FROM MyTable GROUP BY a, c")) + .isInstanceOf(ValidationException.class) + .hasMessageContaining("Please specify an ON CONFLICT clause") + .hasMessageContaining("table.exec.sink.materialized-table-forces-on-conflict-error") + .hasMessageContaining("watermark is missing on") + .hasMessageContaining("MyTable"); + } + @Test void testExplainCreateOrAlterMaterializedTable() { verifyExplain( diff --git a/flink-table/flink-table-planner/src/test/resources/explain/testExplainCreateMaterializedTableDefaultsToOnConflictErrorWhenEnabled.out b/flink-table/flink-table-planner/src/test/resources/explain/testExplainCreateMaterializedTableDefaultsToOnConflictErrorWhenEnabled.out new file mode 100644 index 00000000000000..765b525276c014 --- /dev/null +++ b/flink-table/flink-table-planner/src/test/resources/explain/testExplainCreateMaterializedTableDefaultsToOnConflictErrorWhenEnabled.out @@ -0,0 +1,25 @@ +== Abstract Syntax Tree == +LogicalSink(table=[default_catalog.default_database.MyMTOnConflictTable], fields=[a, cnt]) ++- LogicalProject(a=[$0], cnt=[$2]) + +- LogicalAggregate(group=[{0, 1}], cnt=[COUNT()]) + +- LogicalProject(a=[$0], c=[$1]) + +- LogicalWatermarkAssigner(rowtime=[ts], watermark=[$2]) + +- LogicalTableScan(table=[[default_catalog, default_database, MyWatermarkedTable]]) + +== Optimized Physical Plan == +Sink(table=[default_catalog.default_database.MyMTOnConflictTable], fields=[a, cnt], upsertMaterialize=[true], conflictStrategy=[ERROR]) ++- Calc(select=[a, cnt]) + +- GroupAggregate(groupBy=[a, c], select=[a, c, COUNT(*) AS cnt]) + +- Exchange(distribution=[hash[a, c]]) + +- Calc(select=[a, c]) + +- WatermarkAssigner(rowtime=[ts], watermark=[ts]) + +- TableSourceScan(table=[[default_catalog, default_database, MyWatermarkedTable]], fields=[a, c, ts]) + +== Optimized Execution Plan == +Sink(table=[default_catalog.default_database.MyMTOnConflictTable], fields=[a, cnt], upsertMaterialize=[true], conflictStrategy=[ERROR]) ++- Calc(select=[a, cnt]) + +- GroupAggregate(groupBy=[a, c], select=[a, c, COUNT(*) AS cnt]) + +- Exchange(distribution=[hash[a, c]]) + +- Calc(select=[a, c]) + +- WatermarkAssigner(rowtime=[ts], watermark=[ts]) + +- TableSourceScan(table=[[default_catalog, default_database, MyWatermarkedTable]], fields=[a, c, ts]) diff --git a/flink-table/flink-table-planner/src/test/resources/explain/testExplainCreateMaterializedTableKeepsOldDefaultWhenDisabled.out b/flink-table/flink-table-planner/src/test/resources/explain/testExplainCreateMaterializedTableKeepsOldDefaultWhenDisabled.out new file mode 100644 index 00000000000000..ba18210e8df92b --- /dev/null +++ b/flink-table/flink-table-planner/src/test/resources/explain/testExplainCreateMaterializedTableKeepsOldDefaultWhenDisabled.out @@ -0,0 +1,20 @@ +== Abstract Syntax Tree == +LogicalSink(table=[default_catalog.default_database.MyMTNoConflictDefaultTable], fields=[a, cnt]) ++- LogicalProject(a=[$0], cnt=[$2]) + +- LogicalAggregate(group=[{0, 1}], cnt=[COUNT()]) + +- LogicalProject(a=[$0], c=[$2]) + +- LogicalTableScan(table=[[default_catalog, default_database, MyTable]]) + +== Optimized Physical Plan == +Sink(table=[default_catalog.default_database.MyMTNoConflictDefaultTable], fields=[a, cnt], upsertMaterialize=[true]) ++- Calc(select=[a, cnt]) + +- GroupAggregate(groupBy=[a, c], select=[a, c, COUNT(*) AS cnt]) + +- Exchange(distribution=[hash[a, c]]) + +- TableSourceScan(table=[[default_catalog, default_database, MyTable, project=[a, c], metadata=[]]], fields=[a, c]) + +== Optimized Execution Plan == +Sink(table=[default_catalog.default_database.MyMTNoConflictDefaultTable], fields=[a, cnt], upsertMaterialize=[true]) ++- Calc(select=[a, cnt]) + +- GroupAggregate(groupBy=[a, c], select=[a, c, COUNT(*) AS cnt]) + +- Exchange(distribution=[hash[a, c]]) + +- TableSourceScan(table=[[default_catalog, default_database, MyTable, project=[a, c], metadata=[]]], fields=[a, c]) diff --git a/flink-table/flink-table-planner/src/test/resources/explain/testExplainCreateMaterializedTableKeepsOldDefaultWhenSourceHasNoWatermark.out b/flink-table/flink-table-planner/src/test/resources/explain/testExplainCreateMaterializedTableKeepsOldDefaultWhenSourceHasNoWatermark.out new file mode 100644 index 00000000000000..23cd359a44e216 --- /dev/null +++ b/flink-table/flink-table-planner/src/test/resources/explain/testExplainCreateMaterializedTableKeepsOldDefaultWhenSourceHasNoWatermark.out @@ -0,0 +1,37 @@ +== Abstract Syntax Tree == +LogicalSink(table=[default_catalog.default_database.MyMTJoinedSourcesTable], fields=[a, cnt]) ++- LogicalProject(a=[$0], cnt=[$2]) + +- LogicalAggregate(group=[{0, 1}], cnt=[COUNT()]) + +- LogicalProject(a=[$0], c=[$1]) + +- LogicalJoin(condition=[=($0, $3)], joinType=[inner]) + :- LogicalWatermarkAssigner(rowtime=[ts], watermark=[$2]) + : +- LogicalTableScan(table=[[default_catalog, default_database, MyWatermarkedTable]]) + +- LogicalTableScan(table=[[default_catalog, default_database, MyTable]]) + +== Optimized Physical Plan == +Sink(table=[default_catalog.default_database.MyMTJoinedSourcesTable], fields=[a, cnt], upsertMaterialize=[true]) ++- Calc(select=[a, cnt]) + +- GroupAggregate(groupBy=[a, c], select=[a, c, COUNT(*) AS cnt]) + +- Exchange(distribution=[hash[a, c]]) + +- Calc(select=[a, c]) + +- Join(joinType=[InnerJoin], where=[=(a, a0)], select=[a, c, a0], leftInputSpec=[NoUniqueKey], rightInputSpec=[NoUniqueKey]) + :- Exchange(distribution=[hash[a]]) + : +- Calc(select=[a, c]) + : +- WatermarkAssigner(rowtime=[ts], watermark=[ts]) + : +- TableSourceScan(table=[[default_catalog, default_database, MyWatermarkedTable]], fields=[a, c, ts]) + +- Exchange(distribution=[hash[a]]) + +- TableSourceScan(table=[[default_catalog, default_database, MyTable, project=[a], metadata=[]]], fields=[a]) + +== Optimized Execution Plan == +Sink(table=[default_catalog.default_database.MyMTJoinedSourcesTable], fields=[a, cnt], upsertMaterialize=[true]) ++- Calc(select=[a, cnt]) + +- GroupAggregate(groupBy=[a, c], select=[a, c, COUNT(*) AS cnt]) + +- Exchange(distribution=[hash[a, c]]) + +- Calc(select=[a, c]) + +- Join(joinType=[InnerJoin], where=[(a = a0)], select=[a, c, a0], leftInputSpec=[NoUniqueKey], rightInputSpec=[NoUniqueKey]) + :- Exchange(distribution=[hash[a]]) + : +- Calc(select=[a, c]) + : +- WatermarkAssigner(rowtime=[ts], watermark=[ts]) + : +- TableSourceScan(table=[[default_catalog, default_database, MyWatermarkedTable]], fields=[a, c, ts]) + +- Exchange(distribution=[hash[a]]) + +- TableSourceScan(table=[[default_catalog, default_database, MyTable, project=[a], metadata=[]]], fields=[a]) diff --git a/flink-table/flink-table-planner/src/test/resources/explain/testExplainCreateMaterializedTableNoDuplicatesOnConflictErrorWhenEnabled.out b/flink-table/flink-table-planner/src/test/resources/explain/testExplainCreateMaterializedTableNoDuplicatesOnConflictErrorWhenEnabled.out new file mode 100644 index 00000000000000..2c50b234e98286 --- /dev/null +++ b/flink-table/flink-table-planner/src/test/resources/explain/testExplainCreateMaterializedTableNoDuplicatesOnConflictErrorWhenEnabled.out @@ -0,0 +1,19 @@ +== Abstract Syntax Tree == +LogicalSink(table=[default_catalog.default_database.MyMTNoConflictTable], fields=[a, cnt]) ++- LogicalAggregate(group=[{0}], cnt=[COUNT()]) + +- LogicalProject(a=[COALESCE($0, 0)]) + +- LogicalTableScan(table=[[default_catalog, default_database, MyTable]]) + +== Optimized Physical Plan == +Sink(table=[default_catalog.default_database.MyMTNoConflictTable], fields=[a, cnt]) ++- GroupAggregate(groupBy=[a], select=[a, COUNT(*) AS cnt]) + +- Exchange(distribution=[hash[a]]) + +- Calc(select=[COALESCE(a, 0) AS a]) + +- TableSourceScan(table=[[default_catalog, default_database, MyTable, project=[a], metadata=[]]], fields=[a]) + +== Optimized Execution Plan == +Sink(table=[default_catalog.default_database.MyMTNoConflictTable], fields=[a, cnt]) ++- GroupAggregate(groupBy=[a], select=[a, COUNT(*) AS cnt]) + +- Exchange(distribution=[hash[a]]) + +- Calc(select=[COALESCE(a, 0) AS a]) + +- TableSourceScan(table=[[default_catalog, default_database, MyTable, project=[a], metadata=[]]], fields=[a])