This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-11730-6d9667af63e226d7cec44095128e83d52f62299d in repository https://gitbox.apache.org/repos/asf/seatunnel.git
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 + + ")"); + } }
