+ 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])