This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new 9f243e4eb7 [Fix][Connector-V2] Omit empty UPDATE SET in JDBC MERGE for
all-key tables (#11730)
9f243e4eb7 is described below
commit 9f243e4eb7e56f6e9f71b6844e77d006c85b40c9
Author: yusiwen <[email protected]>
AuthorDate: Fri Oct 2 13:41:58 2026 +0000
[Fix][Connector-V2] Omit empty UPDATE SET in JDBC MERGE for all-key tables
(#11730)
Co-authored-by: yusiwen <[email protected]>
Co-authored-by: zhiweiniu <[email protected]>
---
.../jdbc/internal/dialect/dm/DmdbDialect.java | 13 +++--
.../internal/dialect/oracle/OracleDialect.java | 13 +++--
.../internal/dialect/saphana/SapHanaDialect.java | 15 ++++--
.../dialect/sqlserver/SqlServerDialect.java | 13 +++--
.../internal/dialect/vertica/VerticaDialect.java | 13 +++--
.../jdbc/internal/dialect/xugu/XuguDialect.java | 18 ++++---
.../internal/dialect/yashandb/YashanDbDialect.java | 13 +++--
.../jdbc/internal/dialect/dm/DmdbDialectTest.java | 36 +++++++++++++
.../internal/dialect/oracle/OracleDialectTest.java | 63 ++++++++++++++++++++++
.../dialect/saphana/SapHanaDialectTest.java | 63 ++++++++++++++++++++++
.../dialect/sqlserver/SqlServerDialectTest.java | 38 +++++++++++++
.../dialect/vertica/VerticaDialectTest.java | 58 ++++++++++++++++++++
.../internal/dialect/xugu/XuguDialectTest.java | 63 ++++++++++++++++++++++
.../dialect/yashandb/YashanDbDialectTest.java | 33 ++++++++++++
14 files changed, 426 insertions(+), 26 deletions(-)
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/dm/DmdbDialect.java
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/dm/DmdbDialect.java
index cf9b29baaa..85bd928888 100644
---
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/dm/DmdbDialect.java
+++
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/dm/DmdbDialect.java
@@ -140,19 +140,26 @@ public class DmdbDialect implements JdbcDialect {
// This is compatible with the case that the schema is written or not
written in the conf
// configuration file
String databaseName = tableIdentifier(database, tableName);
+ // When there are no non-unique-key fields to update (e.g. all fields
are unique keys),
+ // the "WHEN MATCHED THEN UPDATE SET" clause must be omitted,
otherwise the database reports
+ // a syntax error because "UPDATE SET" would have an empty body.
+ String matchedClause =
+ StringUtils.isNotBlank(updateSetClause)
+ ? String.format(" WHEN MATCHED THEN UPDATE SET %s ",
updateSetClause)
+ : "";
+
String upsertSQL =
String.format(
" MERGE INTO %s TARGET"
+ " USING (%s) SOURCE"
+ " ON (%s) "
- + " WHEN MATCHED THEN"
- + " UPDATE SET %s"
+ + "%s"
+ " WHEN NOT MATCHED THEN"
+ " INSERT (%s) VALUES (%s)",
databaseName,
usingClause,
onConditions,
- updateSetClause,
+ matchedClause,
insertFields,
insertValues);
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/oracle/OracleDialect.java
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/oracle/OracleDialect.java
index af3f604ea0..2612b2c5cb 100644
---
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/oracle/OracleDialect.java
+++
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/oracle/OracleDialect.java
@@ -177,19 +177,26 @@ public class OracleDialect implements JdbcDialect {
.map(fieldName -> "SOURCE." +
quoteIdentifier(fieldName))
.collect(Collectors.joining(", "));
+ // When there are no non-unique-key fields to update (e.g. all fields
are unique keys),
+ // the "WHEN MATCHED THEN UPDATE SET" clause must be omitted,
otherwise Oracle reports a
+ // syntax error because "UPDATE SET" would have an empty body.
+ String matchedClause =
+ StringUtils.isNotBlank(updateSetClause)
+ ? String.format(" WHEN MATCHED THEN UPDATE SET %s ",
updateSetClause)
+ : "";
+
String upsertSQL =
String.format(
" MERGE INTO %s TARGET"
+ " USING (%s) SOURCE"
+ " ON (%s) "
- + " WHEN MATCHED THEN"
- + " UPDATE SET %s"
+ + "%s"
+ " WHEN NOT MATCHED THEN"
+ " INSERT (%s) VALUES (%s)",
tableIdentifier(database, tableName),
usingClause,
onConditions,
- updateSetClause,
+ matchedClause,
insertFields,
insertValues);
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/saphana/SapHanaDialect.java
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/saphana/SapHanaDialect.java
index 4194d47ff0..874b340191 100644
---
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/saphana/SapHanaDialect.java
+++
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/saphana/SapHanaDialect.java
@@ -18,6 +18,8 @@
package
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.saphana;
+import org.apache.seatunnel.shade.org.apache.commons.lang3.StringUtils;
+
import
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.converter.JdbcRowConverter;
import
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.DatabaseIdentifier;
import
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.JdbcDialect;
@@ -102,19 +104,26 @@ public class SapHanaDialect implements JdbcDialect {
.map(fieldName -> "SOURCE." +
quoteIdentifier(fieldName))
.collect(Collectors.joining(", "));
+ // When there are no non-unique-key fields to update (e.g. all fields
are unique keys),
+ // the "WHEN MATCHED THEN UPDATE SET" clause must be omitted,
otherwise the database reports
+ // a syntax error because "UPDATE SET" would have an empty body.
+ String matchedClause =
+ StringUtils.isNotBlank(updateSetClause)
+ ? String.format(" WHEN MATCHED THEN UPDATE SET %s ",
updateSetClause)
+ : "";
+
String upsertSQL =
String.format(
" MERGE INTO %s AS TARGET"
+ " USING (%s) AS SOURCE"
+ " ON (%s) "
- + " WHEN MATCHED THEN"
- + " UPDATE SET %s"
+ + "%s"
+ " WHEN NOT MATCHED THEN"
+ " INSERT (%s) VALUES (%s)",
tableIdentifier(database, tableName),
usingClause,
onConditions,
- updateSetClause,
+ matchedClause,
insertFields,
insertValues);
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/sqlserver/SqlServerDialect.java
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/sqlserver/SqlServerDialect.java
index 6645ea424c..9639d4f123 100644
---
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/sqlserver/SqlServerDialect.java
+++
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/sqlserver/SqlServerDialect.java
@@ -128,20 +128,27 @@ public class SqlServerDialect implements JdbcDialect {
Arrays.stream(fieldNames)
.map(fieldName -> "[SOURCE]." +
quoteIdentifier(fieldName))
.collect(Collectors.joining(", "));
+ // When there are no non-unique-key fields to update (e.g. all fields
are unique keys),
+ // the "WHEN MATCHED THEN UPDATE SET" clause must be omitted,
otherwise the database reports
+ // a syntax error because "UPDATE SET" would have an empty body.
+ String matchedClause =
+ StringUtils.isNotBlank(updateSetClause)
+ ? String.format(" WHEN MATCHED THEN UPDATE SET %s",
updateSetClause)
+ : "";
+
String upsertSQL =
String.format(
"MERGE INTO %s.%s AS [TARGET]"
+ " USING (%s) AS [SOURCE]"
+ " ON (%s)"
- + " WHEN MATCHED THEN"
- + " UPDATE SET %s"
+ + "%s"
+ " WHEN NOT MATCHED THEN"
+ " INSERT (%s) VALUES (%s);",
quoteDatabaseIdentifier(database),
quoteIdentifier(tableName),
usingClause,
onConditions,
- updateSetClause,
+ matchedClause,
insertFields,
insertValues);
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/vertica/VerticaDialect.java
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/vertica/VerticaDialect.java
index 878ee281a9..35032d9cf2 100644
---
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/vertica/VerticaDialect.java
+++
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/vertica/VerticaDialect.java
@@ -111,20 +111,27 @@ public class VerticaDialect implements JdbcDialect {
.map(fieldName -> "SOURCE." +
quoteIdentifier(fieldName))
.collect(Collectors.joining(", "));
+ // When there are no non-unique-key fields to update (e.g. all fields
are unique keys),
+ // the "WHEN MATCHED THEN UPDATE SET" clause must be omitted,
otherwise the database reports
+ // a syntax error because "UPDATE SET" would have an empty body.
+ String matchedClause =
+ StringUtils.isNotBlank(updateSetClause)
+ ? String.format(" WHEN MATCHED THEN UPDATE SET %s",
updateSetClause)
+ : "";
+
String upsertSQL =
String.format(
" MERGE INTO %s.%s TARGET"
+ " USING (%s) SOURCE"
+ " ON (%s) "
- + " WHEN MATCHED THEN"
- + " UPDATE SET %s"
+ + "%s"
+ " WHEN NOT MATCHED THEN"
+ " INSERT (%s) VALUES (%s)",
quoteDatabaseIdentifier(database),
quoteIdentifier(tableName),
usingClause,
onConditions,
- updateSetClause,
+ matchedClause,
insertFields,
insertValues);
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/xugu/XuguDialect.java
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/xugu/XuguDialect.java
index 19c7763565..a1d518b284 100644
---
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/xugu/XuguDialect.java
+++
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/xugu/XuguDialect.java
@@ -20,7 +20,6 @@ package
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.xugu;
import org.apache.seatunnel.shade.org.apache.commons.lang3.StringUtils;
import org.apache.seatunnel.api.table.catalog.TablePath;
-import org.apache.seatunnel.common.utils.SeaTunnelException;
import
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.converter.JdbcRowConverter;
import
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.DatabaseIdentifier;
import
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.JdbcDialect;
@@ -103,10 +102,6 @@ public class XuguDialect implements JdbcDialect {
Arrays.stream(fieldNames)
.filter(fieldName ->
!Arrays.asList(pkNames).contains(fieldName))
.collect(Collectors.toList());
- if (nonUniqueKeyFields.isEmpty()) {
- throw new SeaTunnelException(
- "The non-primary key field cannot be empty. Please set
other fields");
- }
String valuesBinding =
Arrays.stream(fieldNames)
.map(fieldName -> ":" + fieldName + " " +
quoteIdentifier(fieldName))
@@ -139,19 +134,26 @@ public class XuguDialect implements JdbcDialect {
Arrays.stream(fieldNames)
.map(fieldName -> "SOURCE." +
quoteIdentifier(fieldName))
.collect(Collectors.joining(", "));
+ // When there are no non-unique-key fields to update (e.g. all fields
are unique keys),
+ // the "WHEN MATCHED THEN UPDATE SET" clause must be omitted,
otherwise the database reports
+ // a syntax error because "UPDATE SET" would have an empty body.
+ String matchedClause =
+ StringUtils.isNotBlank(updateSetClause)
+ ? String.format(" WHEN MATCHED THEN UPDATE SET %s ",
updateSetClause)
+ : "";
+
String upsertSQL =
String.format(
" MERGE INTO %s TARGET"
+ " USING (%s) SOURCE"
+ " ON (%s) "
- + " WHEN MATCHED THEN"
- + " UPDATE SET %s"
+ + "%s"
+ " WHEN NOT MATCHED THEN"
+ " INSERT (%s) VALUES (%s)",
tableIdentifier(database, tableName),
usingClause,
onConditions,
- updateSetClause,
+ matchedClause,
insertFields,
insertValues);
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/yashandb/YashanDbDialect.java
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/yashandb/YashanDbDialect.java
index 5d4376402e..a8bd0ef9cf 100644
---
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/yashandb/YashanDbDialect.java
+++
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/yashandb/YashanDbDialect.java
@@ -144,19 +144,26 @@ public class YashanDbDialect implements JdbcDialect {
.map(fieldName -> "SOURCE." +
quoteIdentifier(fieldName))
.collect(Collectors.joining(", "));
+ // When there are no non-unique-key fields to update (e.g. all fields
are unique keys),
+ // the "WHEN MATCHED THEN UPDATE SET" clause must be omitted,
otherwise the database reports
+ // a syntax error because "UPDATE SET" would have an empty body.
+ String matchedClause =
+ StringUtils.isNotBlank(updateSetClause)
+ ? String.format(" WHEN MATCHED THEN UPDATE SET %s ",
updateSetClause)
+ : "";
+
String upsertSQL =
String.format(
" MERGE INTO %s TARGET"
+ " USING (%s) SOURCE"
+ " ON (%s) "
- + " WHEN MATCHED THEN"
- + " UPDATE SET %s"
+ + "%s"
+ " WHEN NOT MATCHED THEN"
+ " INSERT (%s) VALUES (%s)",
tableIdentifier(database, tableName),
usingClause,
onConditions,
- updateSetClause,
+ matchedClause,
insertFields,
insertValues);
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/dm/DmdbDialectTest.java
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/dm/DmdbDialectTest.java
index b1cacdc51b..c04961ee91 100644
---
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/dm/DmdbDialectTest.java
+++
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/dm/DmdbDialectTest.java
@@ -27,6 +27,7 @@ import org.junit.jupiter.api.Test;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
+import java.util.Optional;
public class DmdbDialectTest {
@Test
@@ -159,4 +160,39 @@ public class DmdbDialectTest {
() ->
DmdbDialect.normalizeTablespaceForDdl("MAIN\tTS"));
Assertions.assertTrue(tabCharacter.getMessage().contains("illegal
characters"));
}
+
+ @Test
+ void testAllKeyTableOmitsEmptyUpdateSet() {
+ JdbcDialect dialect = new DmdbDialectFactory().create();
+ String[] allFields = {"id", "name", "age"};
+ Optional<String> upsert =
+ dialect.getUpsertStatement("test_db", "test_table", allFields,
allFields);
+ Assertions.assertTrue(upsert.isPresent(), "upsert statement should be
present");
+ String sql = upsert.get().toUpperCase();
+ Assertions.assertFalse(
+ sql.contains("WHEN MATCHED THEN UPDATE SET"),
+ "all-key table must NOT emit an empty 'WHEN MATCHED THEN
UPDATE SET' (got: "
+ + sql
+ + ")");
+ Assertions.assertTrue(
+ sql.contains("WHEN NOT MATCHED"), "all-key table must still
insert unmatched rows");
+ Assertions.assertTrue(
+ sql.contains("INSERT"), "all-key table statement must contain
an INSERT branch");
+ }
+
+ @Test
+ void testPartialKeyTableStillUpdates() {
+ JdbcDialect dialect = new DmdbDialectFactory().create();
+ String[] allFields = {"id", "name", "age"};
+ String[] uniqueKeys = {"id"};
+ Optional<String> upsert =
+ dialect.getUpsertStatement("test_db", "test_table", allFields,
uniqueKeys);
+ Assertions.assertTrue(upsert.isPresent(), "upsert statement should be
present");
+ String sql = upsert.get().toUpperCase();
+ Assertions.assertTrue(
+ sql.contains("WHEN MATCHED THEN UPDATE SET"),
+ "partial-key table must still emit 'WHEN MATCHED THEN UPDATE
SET' (got: "
+ + sql
+ + ")");
+ }
}
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/oracle/OracleDialectTest.java
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/oracle/OracleDialectTest.java
new file mode 100644
index 0000000000..763646d393
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/oracle/OracleDialectTest.java
@@ -0,0 +1,63 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.oracle;
+
+import
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.JdbcDialect;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.Optional;
+
+public class OracleDialectTest {
+
+ @Test
+ void testAllKeyTableOmitsEmptyUpdateSet() {
+ JdbcDialect dialect = new OracleDialect();
+ String[] allFields = {"id", "name", "age"};
+ Optional<String> upsert =
+ dialect.getUpsertStatement("test_db", "test_table", allFields,
allFields);
+ Assertions.assertTrue(upsert.isPresent(), "upsert statement should be
present");
+ String sql = upsert.get().toUpperCase();
+ Assertions.assertFalse(
+ sql.contains("WHEN MATCHED THEN UPDATE SET"),
+ "all-key table must NOT emit an empty 'WHEN MATCHED THEN
UPDATE SET' (got: "
+ + sql
+ + ")");
+ Assertions.assertTrue(
+ sql.contains("WHEN NOT MATCHED"), "all-key table must still
insert unmatched rows");
+ Assertions.assertTrue(
+ sql.contains("INSERT"), "all-key table statement must contain
an INSERT branch");
+ }
+
+ @Test
+ void testPartialKeyTableStillUpdates() {
+ JdbcDialect dialect = new OracleDialect();
+ String[] allFields = {"id", "name", "age"};
+ String[] uniqueKeys = {"id"};
+ Optional<String> upsert =
+ dialect.getUpsertStatement("test_db", "test_table", allFields,
uniqueKeys);
+ Assertions.assertTrue(upsert.isPresent(), "upsert statement should be
present");
+ String sql = upsert.get().toUpperCase();
+ Assertions.assertTrue(
+ sql.contains("WHEN MATCHED THEN UPDATE SET"),
+ "partial-key table must still emit 'WHEN MATCHED THEN UPDATE
SET' (got: "
+ + sql
+ + ")");
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/saphana/SapHanaDialectTest.java
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/saphana/SapHanaDialectTest.java
new file mode 100644
index 0000000000..480dacfafb
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/saphana/SapHanaDialectTest.java
@@ -0,0 +1,63 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.saphana;
+
+import
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.JdbcDialect;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.Optional;
+
+public class SapHanaDialectTest {
+
+ @Test
+ void testAllKeyTableOmitsEmptyUpdateSet() {
+ JdbcDialect dialect = new SapHanaDialect();
+ String[] allFields = {"id", "name", "age"};
+ Optional<String> upsert =
+ dialect.getUpsertStatement("test_db", "test_table", allFields,
allFields);
+ Assertions.assertTrue(upsert.isPresent(), "upsert statement should be
present");
+ String sql = upsert.get().toUpperCase();
+ Assertions.assertFalse(
+ sql.contains("WHEN MATCHED THEN UPDATE SET"),
+ "all-key table must NOT emit an empty 'WHEN MATCHED THEN
UPDATE SET' (got: "
+ + sql
+ + ")");
+ Assertions.assertTrue(
+ sql.contains("WHEN NOT MATCHED"), "all-key table must still
insert unmatched rows");
+ Assertions.assertTrue(
+ sql.contains("INSERT"), "all-key table statement must contain
an INSERT branch");
+ }
+
+ @Test
+ void testPartialKeyTableStillUpdates() {
+ JdbcDialect dialect = new SapHanaDialect();
+ String[] allFields = {"id", "name", "age"};
+ String[] uniqueKeys = {"id"};
+ Optional<String> upsert =
+ dialect.getUpsertStatement("test_db", "test_table", allFields,
uniqueKeys);
+ Assertions.assertTrue(upsert.isPresent(), "upsert statement should be
present");
+ String sql = upsert.get().toUpperCase();
+ Assertions.assertTrue(
+ sql.contains("WHEN MATCHED THEN UPDATE SET"),
+ "partial-key table must still emit 'WHEN MATCHED THEN UPDATE
SET' (got: "
+ + sql
+ + ")");
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/sqlserver/SqlServerDialectTest.java
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/sqlserver/SqlServerDialectTest.java
index f10d828652..18f5cf1239 100644
---
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/sqlserver/SqlServerDialectTest.java
+++
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/sqlserver/SqlServerDialectTest.java
@@ -21,11 +21,14 @@ import org.apache.seatunnel.api.table.catalog.Column;
import org.apache.seatunnel.api.table.catalog.TableIdentifier;
import org.apache.seatunnel.api.table.catalog.TablePath;
import org.apache.seatunnel.api.table.schema.event.AlterTableChangeColumnEvent;
+import
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.JdbcDialect;
+import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
import java.sql.Connection;
import java.sql.Statement;
+import java.util.Optional;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
@@ -57,4 +60,39 @@ class SqlServerDialectTest {
.execute(
"EXEC [schema_change_test].sys.sp_rename
'dbo.products_sink.add_column2', 'add_column', 'COLUMN';");
}
+
+ @Test
+ void testAllKeyTableOmitsEmptyUpdateSet() {
+ JdbcDialect dialect = new SqlServerDialect();
+ String[] allFields = {"id", "name", "age"};
+ Optional<String> upsert =
+ dialect.getUpsertStatement("test_db", "test_table", allFields,
allFields);
+ Assertions.assertTrue(upsert.isPresent(), "upsert statement should be
present");
+ String sql = upsert.get().toUpperCase();
+ Assertions.assertFalse(
+ sql.contains("WHEN MATCHED THEN UPDATE SET"),
+ "all-key table must NOT emit an empty 'WHEN MATCHED THEN
UPDATE SET' (got: "
+ + sql
+ + ")");
+ Assertions.assertTrue(
+ sql.contains("WHEN NOT MATCHED"), "all-key table must still
insert unmatched rows");
+ Assertions.assertTrue(
+ sql.contains("INSERT"), "all-key table statement must contain
an INSERT branch");
+ }
+
+ @Test
+ void testPartialKeyTableStillUpdates() {
+ JdbcDialect dialect = new SqlServerDialect();
+ String[] allFields = {"id", "name", "age"};
+ String[] uniqueKeys = {"id"};
+ Optional<String> upsert =
+ dialect.getUpsertStatement("test_db", "test_table", allFields,
uniqueKeys);
+ Assertions.assertTrue(upsert.isPresent(), "upsert statement should be
present");
+ String sql = upsert.get().toUpperCase();
+ Assertions.assertTrue(
+ sql.contains("WHEN MATCHED THEN UPDATE SET"),
+ "partial-key table must still emit 'WHEN MATCHED THEN UPDATE
SET' (got: "
+ + sql
+ + ")");
+ }
}
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/vertica/VerticaDialectTest.java
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/vertica/VerticaDialectTest.java
index 1928480f35..eec4b622dd 100644
---
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/vertica/VerticaDialectTest.java
+++
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/vertica/VerticaDialectTest.java
@@ -25,12 +25,14 @@ import org.apache.seatunnel.api.table.catalog.PrimaryKey;
import org.apache.seatunnel.api.table.catalog.TableSchema;
import org.apache.seatunnel.api.table.type.BasicType;
import org.apache.seatunnel.api.table.type.LocalTimeType;
+import
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.JdbcDialect;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
import java.util.Collections;
import java.util.HashMap;
+import java.util.Optional;
public class VerticaDialectTest {
@@ -121,4 +123,60 @@ public class VerticaDialectTest {
upsertCreateTimeSQL,
" MERGE INTO test_database.\"test_table\" TARGET USING (SELECT
CAST(:id AS BIGINT) AS \"id\", CAST(:name AS VARCHAR) AS \"name\", CAST(:age AS
INT) AS \"age\", CAST(:createTime AS TIME) AS \"createTime\" ) SOURCE ON
(TARGET.\"id\"=SOURCE.\"id\" AND TARGET.\"name\"=SOURCE.\"name\" AND
TARGET.\"age\"=SOURCE.\"age\") WHEN MATCHED THEN UPDATE SET
\"createTime\"=SOURCE.\"createTime\" WHEN NOT MATCHED THEN INSERT (\"id\",
\"name\", \"age\", \"createTime\") VALUES (SOURCE.\"id\ [...]
}
+
+ @Test
+ void testAllKeyTableOmitsEmptyUpdateSet() {
+ JdbcDialect dialect = new VerticaDialect();
+ String[] allFields = {"id", "name", "age"};
+ TableSchema tableSchema =
+ TableSchema.builder()
+ .column(
+ PhysicalColumn.of(
+ "id",
+ BasicType.LONG_TYPE,
+ 22L,
+ 0,
+ false,
+ null,
+ "id",
+ "BIGINT",
+ new HashMap<>()))
+ .column(
+ PhysicalColumn.of(
+ "name",
+ BasicType.STRING_TYPE,
+ 128L,
+ 0,
+ false,
+ null,
+ "name",
+ "VARCHAR",
+ new HashMap<>()))
+ .column(
+ PhysicalColumn.of(
+ "age",
+ BasicType.INT_TYPE,
+ null,
+ 0,
+ true,
+ null,
+ "age",
+ "INT",
+ new HashMap<>()))
+ .build();
+ Optional<String> upsert =
+ dialect.getUpsertStatementByTableSchema(
+ "test_db", "test_table", tableSchema, allFields);
+ Assertions.assertTrue(upsert.isPresent(), "upsert statement should be
present");
+ String sql = upsert.get().toUpperCase();
+ Assertions.assertFalse(
+ sql.contains("WHEN MATCHED THEN UPDATE SET"),
+ "all-key table must NOT emit an empty 'WHEN MATCHED THEN
UPDATE SET' (got: "
+ + sql
+ + ")");
+ Assertions.assertTrue(
+ sql.contains("WHEN NOT MATCHED"), "all-key table must still
insert unmatched rows");
+ Assertions.assertTrue(
+ sql.contains("INSERT"), "all-key table statement must contain
an INSERT branch");
+ }
}
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/xugu/XuguDialectTest.java
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/xugu/XuguDialectTest.java
new file mode 100644
index 0000000000..01575504a4
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/xugu/XuguDialectTest.java
@@ -0,0 +1,63 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.xugu;
+
+import
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.JdbcDialect;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.Optional;
+
+public class XuguDialectTest {
+
+ @Test
+ void testAllKeyTableOmitsEmptyUpdateSet() {
+ JdbcDialect dialect = new XuguDialect();
+ String[] allFields = {"id", "name", "age"};
+ Optional<String> upsert =
+ dialect.getUpsertStatement("test_db", "test_table", allFields,
allFields);
+ Assertions.assertTrue(upsert.isPresent(), "upsert statement should be
present");
+ String sql = upsert.get().toUpperCase();
+ Assertions.assertFalse(
+ sql.contains("WHEN MATCHED THEN UPDATE SET"),
+ "all-key table must NOT emit an empty 'WHEN MATCHED THEN
UPDATE SET' (got: "
+ + sql
+ + ")");
+ Assertions.assertTrue(
+ sql.contains("WHEN NOT MATCHED"), "all-key table must still
insert unmatched rows");
+ Assertions.assertTrue(
+ sql.contains("INSERT"), "all-key table statement must contain
an INSERT branch");
+ }
+
+ @Test
+ void testPartialKeyTableStillUpdates() {
+ JdbcDialect dialect = new XuguDialect();
+ String[] allFields = {"id", "name", "age"};
+ String[] uniqueKeys = {"id"};
+ Optional<String> upsert =
+ dialect.getUpsertStatement("test_db", "test_table", allFields,
uniqueKeys);
+ Assertions.assertTrue(upsert.isPresent(), "upsert statement should be
present");
+ String sql = upsert.get().toUpperCase();
+ Assertions.assertTrue(
+ sql.contains("WHEN MATCHED THEN UPDATE SET"),
+ "partial-key table must still emit 'WHEN MATCHED THEN UPDATE
SET' (got: "
+ + sql
+ + ")");
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/yashandb/YashanDbDialectTest.java
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/yashandb/YashanDbDialectTest.java
index 3b96afce2e..c1124e2a27 100644
---
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/yashandb/YashanDbDialectTest.java
+++
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/yashandb/YashanDbDialectTest.java
@@ -224,4 +224,37 @@ public class YashanDbDialectTest {
Assertions.assertTrue(deleteStatement.contains("\"test_table\""));
Assertions.assertTrue(deleteStatement.contains("\"id\""));
}
+
+ @Test
+ public void testAllKeyTableOmitsEmptyUpdateSet() {
+ String[] allFields = {"id", "name", "age"};
+ Optional<String> upsertStatement =
+ DIALECT.getUpsertStatement("test_db", "test_table", allFields,
allFields);
+ Assertions.assertTrue(upsertStatement.isPresent(), "upsert statement
should be present");
+ String sql = upsertStatement.get().toUpperCase();
+ Assertions.assertFalse(
+ sql.contains("WHEN MATCHED THEN UPDATE SET"),
+ "all-key table must NOT emit an empty 'WHEN MATCHED THEN
UPDATE SET' (got: "
+ + sql
+ + ")");
+ Assertions.assertTrue(
+ sql.contains("WHEN NOT MATCHED"), "all-key table must still
insert unmatched rows");
+ Assertions.assertTrue(
+ sql.contains("INSERT"), "all-key table statement must contain
an INSERT branch");
+ }
+
+ @Test
+ public void testPartialKeyTableStillUpdates() {
+ String[] allFields = {"id", "name", "age"};
+ String[] uniqueKeyFields = {"id"};
+ Optional<String> upsertStatement =
+ DIALECT.getUpsertStatement("test_db", "test_table", allFields,
uniqueKeyFields);
+ Assertions.assertTrue(upsertStatement.isPresent(), "upsert statement
should be present");
+ String sql = upsertStatement.get().toUpperCase();
+ Assertions.assertTrue(
+ sql.contains("WHEN MATCHED THEN UPDATE SET"),
+ "partial-key table must still emit 'WHEN MATCHED THEN UPDATE
SET' (got: "
+ + sql
+ + ")");
+ }
}