This is an automated email from the ASF dual-hosted git repository.

davidzollo 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 b60eb4c562 [Feature][Connector-V2] PR1: Pass sink table-options into 
auto-created MySQL target tables (#11101)
b60eb4c562 is described below

commit b60eb4c56258babc91b1b9149a392be65c253718
Author: luxiaolong <[email protected]>
AuthorDate: Wed Jul 1 11:30:28 2026 +0800

    [Feature][Connector-V2] PR1: Pass sink table-options into auto-created 
MySQL target tables (#11101)
    
    Co-authored-by: luxiaolong-ct 
<[email protected]>
    Co-authored-by: det101 <[email protected]>
---
 docs/en/connectors/sink/Jdbc.md                    |  39 ++++
 docs/zh/connectors/sink/Jdbc.md                    |  39 ++++
 .../api/options/SinkConnectorCommonOptions.java    |  14 ++
 .../jdbc/internal/dialect/JdbcDialect.java         |  19 ++
 .../jdbc/internal/dialect/mysql/MysqlDialect.java  |  27 +++
 .../seatunnel/jdbc/sink/JdbcSinkFactory.java       |   7 +
 .../sink/JdbcTableOptionsConditionExtension.java   |  57 ++++++
 .../jdbc/sink/JdbcTableOptionsValidator.java       |  57 ++++++
 .../mysql/MysqlCreateTableSqlBuilderTest.java      |  32 ++++
 .../internal/dialect/mysql/MysqlDialectTest.java   |  69 +++++++
 .../JdbcTableOptionsConditionExtensionTest.java    |  97 ++++++++++
 .../seatunnel/jdbc/JdbcMysqlTableOptionsIT.java    | 203 +++++++++++++++++++++
 .../jdbc_mysql_sink_with_table_options.conf        |  54 ++++++
 13 files changed, 714 insertions(+)

diff --git a/docs/en/connectors/sink/Jdbc.md b/docs/en/connectors/sink/Jdbc.md
index c5ee8ccd79..dfb1b77a43 100644
--- a/docs/en/connectors/sink/Jdbc.md
+++ b/docs/en/connectors/sink/Jdbc.md
@@ -60,6 +60,7 @@ support `Xa transactions`. You can set `is_exactly_once=true` 
to enable it.
 | data_save_mode                            | Enum    | No       | APPEND_DATA 
                 |
 | custom_sql                                | String  | No       | -           
                 |
 | enable_upsert                             | Boolean | No       | true        
                 |
+| table_options                             | Map     | No       | -           
                 |
 | use_copy_statement                        | Boolean | No       | false       
                 |
 | oracle_insert_mode                        | Enum    | No       | 
CONVENTIONAL                 |
 | create_index                              | Boolean | No       | true        
                 |
@@ -230,6 +231,44 @@ When data_save_mode selects CUSTOM_PROCESSING, you should 
fill in the CUSTOM_SQL
 
 Note: in sink `query` mode, `custom_sql` is not executed. This behavior is a 
current limitation of JDBC sink.
 
+### table_options [Map]
+
+Sink-specific table options applied when SaveMode creates the target table 
(DDL phase). They take effect only when `schema_save_mode` triggers table 
creation, such as `CREATE_SCHEMA_WHEN_NOT_EXIST` or `RECREATE_SCHEMA`. They do 
**not** affect INSERT/UPSERT at runtime and do **not** run `ALTER TABLE` on 
existing tables.
+
+Current support:
+
+| Dialect | Supported | Allowed keys |
+|---------|-----------|--------------|
+| MySQL | Yes | `engine`, `charset`, `collate` |
+| Other JDBC dialects | No | Non-empty `table_options` fails validation at job 
submission |
+
+Invalid or unsupported keys are validated early via `JdbcSinkFactory` option 
rules (`--check` and job submission), not only at runtime DDL.
+
+Example (MySQL auto-create with engine and charset):
+
+```hocon
+sink {
+  Jdbc {
+    url = "jdbc:mysql://localhost:3307/mydb"
+    driver = "com.mysql.cj.jdbc.Driver"
+    username = "root"
+    password = "password"
+    database = "mydb"
+    table = "orders"
+    generate_sink_sql = true
+    schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST"
+    primary_keys = ["id"]
+    table_options = {
+      "engine" = "InnoDB"
+      "charset" = "utf8mb4"
+      "collate" = "utf8mb4_general_ci"
+    }
+  }
+}
+```
+
+The generated `CREATE TABLE` statement appends `ENGINE`, `DEFAULT CHARSET`, 
and `COLLATE` clauses. Keys outside the dialect whitelist (for example 
`bucket_num`) fail during job submission.
+
 ### enable_upsert [boolean]
 
 Enable upsert by primary_keys exist, If the task has no key duplicate data, 
setting this parameter to `false` can speed up data import
diff --git a/docs/zh/connectors/sink/Jdbc.md b/docs/zh/connectors/sink/Jdbc.md
index 732eaa8df6..17783d2eaa 100644
--- a/docs/zh/connectors/sink/Jdbc.md
+++ b/docs/zh/connectors/sink/Jdbc.md
@@ -58,6 +58,7 @@ import ChangeLog from '../changelog/connector-jdbc.md';
 | data_save_mode                            | Enum    | 否    | APPEND_DATA     
             |
 | custom_sql                                | String  | 否    | -               
             |
 | enable_upsert                             | Boolean | 否    | true            
             |
+| table_options                             | Map     | 否    | -               
             |
 | use_copy_statement                        | Boolean | 否    | false           
             |
 | oracle_insert_mode                        | Enum    | 否    | CONVENTIONAL    
             |
 | access_key_id                             | String  | 否       |              
                |
@@ -218,6 +219,44 @@ Sink插件常用参数,请参考 [Sink常用选项](../common-options/sink-com
 `CUSTOM_PROCESSING`:允许用户自定义数据处理方式<br/>
 `ERROR_WHEN_DATA_EXISTS`:当有数据时抛出错误<br/>
 
+### table_options [Map]
+
+Sink 在自动建表(SaveMode DDL)时附加的表级选项。仅在 `schema_save_mode` 触发建表时生效,例如 
`CREATE_SCHEMA_WHEN_NOT_EXIST`、`RECREATE_SCHEMA`;**不影响**数据写入阶段的 
INSERT/UPSERT,也**不会**对已存在表执行 `ALTER TABLE`。
+
+当前支持情况:
+
+| 方言 | 是否支持 | 可用 key |
+|------|----------|----------|
+| MySQL | 是 | `engine`、`charset`、`collate` |
+| 其他 JDBC 方言 | 否 | 配置非空 `table_options` 时任务启动即校验失败 |
+
+非法或不支持的 key 会在 `JdbcSinkFactory` 的 option 规则阶段提前校验(`--check` 与作业提交),而非仅在运行时 
DDL 阶段失败。
+
+示例(MySQL 自动建表时指定存储引擎与字符集):
+
+```hocon
+sink {
+  Jdbc {
+    url = "jdbc:mysql://localhost:3307/mydb"
+    driver = "com.mysql.cj.jdbc.Driver"
+    username = "root"
+    password = "password"
+    database = "mydb"
+    table = "orders"
+    generate_sink_sql = true
+    schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST"
+    primary_keys = ["id"]
+    table_options = {
+      "engine" = "InnoDB"
+      "charset" = "utf8mb4"
+      "collate" = "utf8mb4_general_ci"
+    }
+  }
+}
+```
+
+生成的 DDL 会追加 `ENGINE`、`DEFAULT CHARSET`、`COLLATE` 子句。未在白名单内的 key(如 
`bucket_num`)会在作业提交阶段报错。
+
 ### custom_sql [String]
 
 
当`data_save_mode`选择`CUSTOM_PROCESSING`时,需要填写`CUSTOM_SQL`参数。该参数通常填写一条可以执行的SQL。SQL将在同步任务之前执行
diff --git 
a/seatunnel-api/src/main/java/org/apache/seatunnel/api/options/SinkConnectorCommonOptions.java
 
b/seatunnel-api/src/main/java/org/apache/seatunnel/api/options/SinkConnectorCommonOptions.java
index c245d908dc..1808aad506 100644
--- 
a/seatunnel-api/src/main/java/org/apache/seatunnel/api/options/SinkConnectorCommonOptions.java
+++ 
b/seatunnel-api/src/main/java/org/apache/seatunnel/api/options/SinkConnectorCommonOptions.java
@@ -21,6 +21,9 @@ import org.apache.seatunnel.api.annotation.Experimental;
 import org.apache.seatunnel.api.configuration.Option;
 import org.apache.seatunnel.api.configuration.Options;
 
+import java.util.HashMap;
+import java.util.Map;
+
 public class SinkConnectorCommonOptions extends ConnectorCommonOptions {
 
     @Experimental
@@ -29,4 +32,15 @@ public class SinkConnectorCommonOptions extends 
ConnectorCommonOptions {
                     .intType()
                     .defaultValue(1)
                     .withDescription("The replica number of multi table sink 
writer");
+
+    @Experimental
+    public static Option<Map<String, String>> TABLE_OPTIONS =
+            Options.key("table_options")
+                    .mapType()
+                    .defaultValue(new HashMap<>())
+                    .withDescription(
+                            "Experimental sink-specific table options applied 
when auto-creating "
+                                    + "target tables during SaveMode. Allowed 
keys and semantics are "
+                                    + "defined and validated by each sink 
connector and/or database "
+                                    + "dialect; see the connector 
documentation for supported options.");
 }
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/JdbcDialect.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/JdbcDialect.java
index d943e22d8a..95055a7d70 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/JdbcDialect.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/JdbcDialect.java
@@ -19,6 +19,7 @@ package 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect;
 
 import org.apache.seatunnel.shade.org.apache.commons.lang3.StringUtils;
 
+import org.apache.seatunnel.api.common.SeaTunnelAPIErrorCode;
 import org.apache.seatunnel.api.table.catalog.TablePath;
 import org.apache.seatunnel.api.table.catalog.TableSchema;
 import org.apache.seatunnel.api.table.converter.BasicTypeDefine;
@@ -32,6 +33,7 @@ import 
org.apache.seatunnel.api.table.schema.event.AlterTableModifyColumnEvent;
 import org.apache.seatunnel.api.table.schema.event.SchemaChangeEvent;
 import org.apache.seatunnel.api.table.type.SqlType;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.config.JdbcConnectionConfig;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.exception.JdbcConnectorException;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.connection.JdbcConnectionProvider;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.connection.SimpleJdbcConnectionProvider;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.converter.JdbcRowConverter;
@@ -904,4 +906,21 @@ public interface JdbcDialect extends Serializable {
     default String dualTable() {
         return "";
     }
+
+    /**
+     * Validate sink table options for auto-create mode.
+     *
+     * <p>Default behavior is fail-fast for any non-empty table options. 
Dialects should override
+     * this when they support sink-specific table options.
+     */
+    default void validateTableOptions(Map<String, String> tableOptions) {
+        if (tableOptions == null || tableOptions.isEmpty()) {
+            return;
+        }
+        throw new JdbcConnectorException(
+                SeaTunnelAPIErrorCode.CONFIG_VALIDATION_FAILED,
+                String.format(
+                        "JDBC table_options are not supported for dialect '%s' 
yet.",
+                        dialectName()));
+    }
 }
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/mysql/MysqlDialect.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/mysql/MysqlDialect.java
index c63cbc4097..09fefdd564 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/mysql/MysqlDialect.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/mysql/MysqlDialect.java
@@ -19,9 +19,11 @@ package 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.mysql;
 
 import org.apache.seatunnel.shade.org.apache.commons.lang3.StringUtils;
 
+import org.apache.seatunnel.api.common.SeaTunnelAPIErrorCode;
 import org.apache.seatunnel.api.table.catalog.TablePath;
 import org.apache.seatunnel.api.table.converter.BasicTypeDefine;
 import org.apache.seatunnel.api.table.converter.TypeConverter;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.exception.JdbcConnectorException;
 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;
@@ -41,11 +43,14 @@ import java.sql.SQLException;
 import java.sql.Statement;
 import java.util.ArrayList;
 import java.util.Arrays;
+import java.util.Collections;
 import java.util.HashMap;
+import java.util.LinkedHashSet;
 import java.util.List;
 import java.util.Locale;
 import java.util.Map;
 import java.util.Optional;
+import java.util.Set;
 import java.util.stream.Collectors;
 
 @Slf4j
@@ -53,6 +58,9 @@ public class MysqlDialect implements JdbcDialect {
 
     private static final List NOT_SUPPORTED_DEFAULT_VALUES =
             Arrays.asList(MysqlType.BLOB, MysqlType.TEXT, MysqlType.JSON, 
MysqlType.GEOMETRY);
+    private static final Set<String> SUPPORTED_TABLE_OPTIONS =
+            Collections.unmodifiableSet(
+                    new LinkedHashSet<>(Arrays.asList("engine", "charset", 
"collate")));
 
     public String fieldIde = FieldIdeEnum.ORIGINAL.getValue();
 
@@ -389,4 +397,23 @@ public class MysqlDialect implements JdbcDialect {
                 return false;
         }
     }
+
+    @Override
+    public void validateTableOptions(Map<String, String> tableOptions) {
+        if (tableOptions == null || tableOptions.isEmpty()) {
+            return;
+        }
+
+        Set<String> unsupportedOptions = new 
LinkedHashSet<>(tableOptions.keySet());
+        unsupportedOptions.removeAll(SUPPORTED_TABLE_OPTIONS);
+        if (!unsupportedOptions.isEmpty()) {
+            throw new JdbcConnectorException(
+                    SeaTunnelAPIErrorCode.CONFIG_VALIDATION_FAILED,
+                    String.format(
+                            "Unsupported JDBC table_options for dialect '%s': 
%s. Supported keys: %s",
+                            dialectName(),
+                            String.join(", ", unsupportedOptions),
+                            String.join(", ", SUPPORTED_TABLE_OPTIONS)));
+        }
+    }
 }
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcSinkFactory.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcSinkFactory.java
index bf035c2c0b..33cb078fc1 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcSinkFactory.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcSinkFactory.java
@@ -73,6 +73,7 @@ public class JdbcSinkFactory implements TableSinkFactory {
     @Override
     public TableSink createSink(TableSinkFactoryContext context) {
         ReadonlyConfig config = context.getOptions();
+        Map<String, String> sinkTableOptions = 
config.get(SinkConnectorCommonOptions.TABLE_OPTIONS);
         CatalogTable catalogTable = context.getCatalogTable();
         ReadonlyConfig catalogOptions = getCatalogOptions(context);
         Optional<String> optionalTable = 
config.getOptional(JdbcSinkOptions.TABLE);
@@ -182,6 +183,7 @@ public class JdbcSinkFactory implements TableSinkFactory {
         final ReadonlyConfig options = config;
         JdbcSinkConfig sinkConfig = JdbcSinkConfig.of(config);
         FieldIdeEnum fieldIdeEnum = config.get(JdbcSinkOptions.FIELD_IDE);
+        catalogTable.getOptions().putAll(sinkTableOptions);
         catalogTable
                 .getOptions()
                 .put("fieldIde", fieldIdeEnum == null ? null : 
fieldIdeEnum.getValue());
@@ -246,6 +248,11 @@ public class JdbcSinkFactory implements TableSinkFactory {
                         JdbcSinkOptions.TABLE_SUFFIX,
                         SinkConnectorCommonOptions.MULTI_TABLE_SINK_REPLICA,
                         JdbcSinkOptions.DIALECT)
+                .optional(
+                        SinkConnectorCommonOptions.TABLE_OPTIONS,
+                        Conditions.extension(
+                                SinkConnectorCommonOptions.TABLE_OPTIONS,
+                                JdbcTableOptionsConditionExtension.INSTANCE))
                 .conditional(
                         JdbcSinkOptions.IS_EXACTLY_ONCE,
                         true,
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcTableOptionsConditionExtension.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcTableOptionsConditionExtension.java
new file mode 100644
index 0000000000..dc74d956c8
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcTableOptionsConditionExtension.java
@@ -0,0 +1,57 @@
+/*
+ * 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.sink;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConditionExtension;
+import org.apache.seatunnel.api.configuration.util.OptionValidationException;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.exception.JdbcConnectorException;
+
+import java.util.Map;
+
+/**
+ * Early validation for JDBC sink {@code table_options}. Delegates to {@link
+ * JdbcTableOptionsValidator} so dialect-specific rules are defined on {@link
+ * 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.JdbcDialect}.
+ */
+public class JdbcTableOptionsConditionExtension implements 
ConditionExtension<Map<String, String>> {
+
+    public static final JdbcTableOptionsConditionExtension INSTANCE =
+            new JdbcTableOptionsConditionExtension();
+
+    private JdbcTableOptionsConditionExtension() {}
+
+    @Override
+    public String description() {
+        return "must use dialect-specific keys supported by the JDBC sink (see 
JDBC connector docs)";
+    }
+
+    @Override
+    public boolean evaluate(ReadonlyConfig config, Map<String, String> value)
+            throws OptionValidationException {
+        if (value == null || value.isEmpty()) {
+            return true;
+        }
+        try {
+            JdbcTableOptionsValidator.validate(config, value);
+            return true;
+        } catch (JdbcConnectorException e) {
+            throw new OptionValidationException(e.getMessage());
+        }
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcTableOptionsValidator.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcTableOptionsValidator.java
new file mode 100644
index 0000000000..00d6b69c84
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcTableOptionsValidator.java
@@ -0,0 +1,57 @@
+/*
+ * 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.sink;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.options.SinkConnectorCommonOptions;
+import org.apache.seatunnel.connectors.seatunnel.jdbc.config.JdbcSinkConfig;
+import org.apache.seatunnel.connectors.seatunnel.jdbc.config.JdbcSinkOptions;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.JdbcDialect;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.JdbcDialectLoader;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.dialectenum.FieldIdeEnum;
+
+import java.util.Collections;
+import java.util.Map;
+
+/** Validates sink {@code table_options} via the resolved {@link JdbcDialect}. 
*/
+public final class JdbcTableOptionsValidator {
+
+    private JdbcTableOptionsValidator() {}
+
+    public static void validate(ReadonlyConfig config, Map<String, String> 
tableOptions) {
+        if (tableOptions == null || tableOptions.isEmpty()) {
+            return;
+        }
+        JdbcSinkConfig sinkConfig = JdbcSinkConfig.of(config);
+        FieldIdeEnum fieldIdeEnum = config.get(JdbcSinkOptions.FIELD_IDE);
+        JdbcDialect dialect =
+                JdbcDialectLoader.load(
+                        sinkConfig.getJdbcConnectionConfig().getUrl(),
+                        
sinkConfig.getJdbcConnectionConfig().getCompatibleMode(),
+                        sinkConfig.getJdbcConnectionConfig().getDialect(),
+                        fieldIdeEnum == null ? null : fieldIdeEnum.getValue());
+        dialect.validateTableOptions(tableOptions);
+    }
+
+    public static void validate(ReadonlyConfig config) {
+        validate(
+                config,
+                config.getOptional(SinkConnectorCommonOptions.TABLE_OPTIONS)
+                        .orElse(Collections.emptyMap()));
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/mysql/MysqlCreateTableSqlBuilderTest.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/mysql/MysqlCreateTableSqlBuilderTest.java
index 0f62186e20..21740a282f 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/mysql/MysqlCreateTableSqlBuilderTest.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/mysql/MysqlCreateTableSqlBuilderTest.java
@@ -40,7 +40,9 @@ import org.junit.jupiter.api.Test;
 import java.io.PrintStream;
 import java.util.ArrayList;
 import java.util.Arrays;
+import java.util.Collections;
 import java.util.HashMap;
+import java.util.Map;
 
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.when;
@@ -150,6 +152,36 @@ public class MysqlCreateTableSqlBuilderTest {
         Assertions.assertEquals(expectSkipIndex, createTableSqlSkipIndex);
     }
 
+    @Test
+    public void testBuildCreateTableSqlWithTableOptions() {
+        TablePath tablePath = TablePath.of("test_db", "test_table");
+        TableSchema tableSchema =
+                TableSchema.builder()
+                        .column(PhysicalColumn.of("id", BasicType.LONG_TYPE, 
0, false, null, "id"))
+                        .primaryKey(PrimaryKey.of("id", 
Lists.newArrayList("id")))
+                        .build();
+        Map<String, String> options = new HashMap<>();
+        options.put(MySqlCatalog.TABLE_OPTION_ENGINE, "InnoDB");
+        options.put(MySqlCatalog.TABLE_OPTION_CHARSET, "utf8mb4");
+        options.put(MySqlCatalog.TABLE_OPTION_COLLATE, "utf8mb4_unicode_ci");
+        CatalogTable catalogTable =
+                CatalogTable.of(
+                        TableIdentifier.of("test_catalog", "test_db", 
"test_table"),
+                        tableSchema,
+                        options,
+                        Collections.emptyList(),
+                        "table with options");
+
+        String createTableSql =
+                MysqlCreateTableSqlBuilder.builder(
+                                tablePath, catalogTable, 
MySqlTypeConverter.DEFAULT_INSTANCE, true)
+                        .build(DatabaseIdentifier.MYSQL);
+
+        Assertions.assertTrue(createTableSql.contains("ENGINE = InnoDB"));
+        Assertions.assertTrue(createTableSql.contains("DEFAULT CHARSET = 
utf8mb4"));
+        Assertions.assertTrue(createTableSql.contains("COLLATE = 
utf8mb4_unicode_ci"));
+    }
+
     @Test
     public void testColumnSinkType() {
         MysqlCreateTableSqlBuilder sqlBuilder = 
mock(MysqlCreateTableSqlBuilder.class);
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/mysql/MysqlDialectTest.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/mysql/MysqlDialectTest.java
index 122cec9e28..1c9925eb9c 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/mysql/MysqlDialectTest.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/mysql/MysqlDialectTest.java
@@ -18,6 +18,11 @@
 package org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.mysql;
 
 import org.apache.seatunnel.api.table.catalog.TablePath;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.exception.JdbcConnectorException;
+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;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.JdbcDialectTypeMapper;
 import org.apache.seatunnel.connectors.seatunnel.jdbc.source.JdbcSourceTable;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.source.StringRangeSplitDecision;
 
@@ -33,15 +38,79 @@ import java.sql.ResultSet;
 import java.sql.Statement;
 import java.util.ArrayList;
 import java.util.Arrays;
+import java.util.Collections;
 import java.util.HashMap;
 import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
+import java.util.Optional;
 import java.util.zip.CRC32;
 
 @Slf4j
 public class MysqlDialectTest {
 
+    @Test
+    public void testValidateTableOptionsForMysql() {
+        MysqlDialect dialect = new MysqlDialect();
+        Map<String, String> tableOptions = new HashMap<>();
+        tableOptions.put("engine", "InnoDB");
+        tableOptions.put("charset", "utf8mb4");
+        tableOptions.put("collate", "utf8mb4_unicode_ci");
+
+        Assertions.assertDoesNotThrow(() -> 
dialect.validateTableOptions(tableOptions));
+    }
+
+    @Test
+    public void testValidateTableOptionsForMysqlWithUnknownOption() {
+        MysqlDialect dialect = new MysqlDialect();
+        Map<String, String> tableOptions = new HashMap<>();
+        tableOptions.put("bucket_num", "3");
+
+        JdbcConnectorException exception =
+                Assertions.assertThrows(
+                        JdbcConnectorException.class,
+                        () -> dialect.validateTableOptions(tableOptions));
+        Assertions.assertTrue(exception.getMessage().contains("Unsupported 
JDBC table_options"));
+    }
+
+    @Test
+    public void testValidateTableOptionsForUnsupportedDialect() {
+        JdbcDialect unsupportedDialect =
+                new JdbcDialect() {
+                    @Override
+                    public String dialectName() {
+                        return DatabaseIdentifier.POSTGRESQL;
+                    }
+
+                    @Override
+                    public JdbcRowConverter getRowConverter() {
+                        return null;
+                    }
+
+                    @Override
+                    public JdbcDialectTypeMapper getJdbcDialectTypeMapper() {
+                        return null;
+                    }
+
+                    @Override
+                    public Optional<String> getUpsertStatement(
+                            String database,
+                            String tableName,
+                            String[] fieldNames,
+                            String[] pkNames) {
+                        return Optional.empty();
+                    }
+                };
+
+        JdbcConnectorException exception =
+                Assertions.assertThrows(
+                        JdbcConnectorException.class,
+                        () ->
+                                unsupportedDialect.validateTableOptions(
+                                        Collections.singletonMap("engine", 
"InnoDB")));
+        Assertions.assertTrue(exception.getMessage().contains("not 
supported"));
+    }
+
     @Test
     public void testValidateStringRangeSplitAcceptsPrintableAsciiPunctuation() 
throws Exception {
         MysqlDialect dialect = new MysqlDialect();
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcTableOptionsConditionExtensionTest.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcTableOptionsConditionExtensionTest.java
new file mode 100644
index 0000000000..dd8e99697c
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcTableOptionsConditionExtensionTest.java
@@ -0,0 +1,97 @@
+/*
+ * 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.sink;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConfigValidator;
+import org.apache.seatunnel.api.configuration.util.OptionValidationException;
+import org.apache.seatunnel.api.options.SinkConnectorCommonOptions;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.HashMap;
+import java.util.Map;
+
+/**
+ * Verifies {@code table_options} early validation is wired through {@link
+ * JdbcSinkFactory#optionRule()} and {@link 
JdbcTableOptionsConditionExtension}. Dialect-specific
+ * allowlists are covered in {@code *DialectTest} classes.
+ */
+class JdbcTableOptionsConditionExtensionTest {
+
+    @Test
+    void testMysqlTableOptionsPassViaOptionRule() {
+        Map<String, Object> config = mysqlSinkConfig();
+        Map<String, String> tableOptions = new HashMap<>();
+        tableOptions.put("engine", "InnoDB");
+        tableOptions.put("charset", "utf8mb4");
+        tableOptions.put("collate", "utf8mb4_unicode_ci");
+        config.put(SinkConnectorCommonOptions.TABLE_OPTIONS.key(), 
tableOptions);
+
+        Assertions.assertDoesNotThrow(() -> validateSinkOptionRule(config));
+    }
+
+    @Test
+    void testPostgresRejectsNonEmptyTableOptionsViaOptionRule() {
+        Map<String, Object> config = postgresSinkConfig();
+        Map<String, String> tableOptions = new HashMap<>();
+        tableOptions.put("fillfactor", "70");
+        config.put(SinkConnectorCommonOptions.TABLE_OPTIONS.key(), 
tableOptions);
+
+        OptionValidationException exception =
+                Assertions.assertThrows(
+                        OptionValidationException.class, () -> 
validateSinkOptionRule(config));
+        Assertions.assertTrue(
+                exception.getMessage().contains("not supported for dialect 
'Postgres'"));
+    }
+
+    @Test
+    void testAbsentTableOptionsSkipsExtension() {
+        Assertions.assertDoesNotThrow(() -> 
validateSinkOptionRule(mysqlSinkConfig()));
+    }
+
+    @Test
+    void testEmptyTableOptionsSkipsExtension() {
+        Map<String, Object> config = mysqlSinkConfig();
+        config.put(SinkConnectorCommonOptions.TABLE_OPTIONS.key(), new 
HashMap<>());
+
+        Assertions.assertDoesNotThrow(() -> validateSinkOptionRule(config));
+    }
+
+    private static void validateSinkOptionRule(Map<String, Object> config) {
+        ConfigValidator.of(ReadonlyConfig.fromMap(config))
+                .validate(new JdbcSinkFactory().optionRule());
+    }
+
+    private static Map<String, Object> mysqlSinkConfig() {
+        Map<String, Object> config = new HashMap<>();
+        config.put("url", "jdbc:mysql://127.0.0.1:3306/test");
+        config.put("driver", "com.mysql.cj.jdbc.Driver");
+        config.put("query", "INSERT INTO test_table VALUES (?)");
+        return config;
+    }
+
+    private static Map<String, Object> postgresSinkConfig() {
+        Map<String, Object> config = new HashMap<>();
+        config.put("url", "jdbc:postgresql://127.0.0.1:5432/test");
+        config.put("driver", "org.postgresql.Driver");
+        config.put("query", "INSERT INTO test_table VALUES (?)");
+        return config;
+    }
+}
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-7/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/JdbcMysqlTableOptionsIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-7/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/JdbcMysqlTableOptionsIT.java
new file mode 100644
index 0000000000..0d3a755fec
--- /dev/null
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-7/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/JdbcMysqlTableOptionsIT.java
@@ -0,0 +1,203 @@
+/*
+ * 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;
+
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.catalog.mysql.MySqlCatalog;
+import org.apache.seatunnel.e2e.common.TestResource;
+import org.apache.seatunnel.e2e.common.TestSuiteBase;
+import org.apache.seatunnel.e2e.common.container.ContainerExtendedFactory;
+import org.apache.seatunnel.e2e.common.container.TestContainer;
+import org.apache.seatunnel.e2e.common.junit.TestContainerExtension;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.TestTemplate;
+import org.testcontainers.containers.Container;
+import org.testcontainers.containers.MySQLContainer;
+import org.testcontainers.containers.output.Slf4jLogConsumer;
+import org.testcontainers.containers.wait.strategy.Wait;
+import org.testcontainers.images.PullPolicy;
+import org.testcontainers.lifecycle.Startables;
+import org.testcontainers.utility.DockerImageName;
+import org.testcontainers.utility.DockerLoggerFactory;
+
+import lombok.extern.slf4j.Slf4j;
+
+import java.io.IOException;
+import java.sql.Connection;
+import java.sql.DriverManager;
+import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.sql.Statement;
+import java.time.Duration;
+import java.util.stream.Stream;
+
+@Slf4j
+public class JdbcMysqlTableOptionsIT extends TestSuiteBase implements 
TestResource {
+
+    private static final String MYSQL_DRIVER_JAR =
+            
"https://repo1.maven.org/maven2/com/mysql/mysql-connector-j/8.0.32/mysql-connector-j-8.0.32.jar";;
+
+    private static final String MYSQL_IMAGE = "mysql:8.0.43";
+    private static final String MYSQL_CONTAINER_HOST = 
"mysql-e2e-table-options";
+    private static final String MYSQL_DATABASE = "seatunnel";
+    private static final String MYSQL_SOURCE = "source";
+    private static final String MYSQL_SINK = "sink_table_options";
+
+    private static final String MYSQL_USERNAME = "root";
+    private static final String MYSQL_PASSWORD = "Abc!@#135_seatunnel";
+
+    private static final String CONFIG_FILE = 
"/jdbc_mysql_sink_with_table_options.conf";
+
+    private static final String CREATE_SOURCE_TABLE_SQL =
+            "CREATE TABLE IF NOT EXISTS `"
+                    + MYSQL_SOURCE
+                    + "` (\n"
+                    + "    `id`   BIGINT       NOT NULL,\n"
+                    + "    `name` VARCHAR(255) DEFAULT NULL,\n"
+                    + "    PRIMARY KEY (`id`)\n"
+                    + ");";
+
+    private static final String INSERT_SOURCE_SQL =
+            "INSERT INTO `"
+                    + MYSQL_SOURCE
+                    + "` (`id`, `name`) VALUES (1, 'name_1'), (2, 'name_2'), 
(3, 'name_3');";
+
+    // MySQL 8.0.43 cold start may exceed the default 120s JDBC wait in 
Testcontainers.
+    private static final int MYSQL_STARTUP_TIMEOUT_SECONDS =
+            (int) Duration.ofMinutes(10).getSeconds();
+
+    private MySQLContainer<?> mysqlContainer;
+
+    @TestContainerExtension
+    protected final ContainerExtendedFactory extendedFactory =
+            container -> {
+                Container.ExecResult extraCommands =
+                        container.execInContainer(
+                                "bash",
+                                "-c",
+                                "mkdir -p /tmp/seatunnel/plugins/Jdbc/lib && 
cd /tmp/seatunnel/plugins/Jdbc/lib && wget "
+                                        + MYSQL_DRIVER_JAR);
+                Assertions.assertEquals(0, extraCommands.getExitCode(), 
extraCommands.getStderr());
+            };
+
+    void initContainer() {
+        DockerImageName imageName = DockerImageName.parse(MYSQL_IMAGE);
+        mysqlContainer =
+                new MySQLContainer<>(imageName)
+                        
.withImagePullPolicy(PullPolicy.ageBased(Duration.ofDays(7)))
+                        .withUsername(MYSQL_USERNAME)
+                        .withPassword(MYSQL_PASSWORD)
+                        .withDatabaseName(MYSQL_DATABASE)
+                        .withNetwork(NETWORK)
+                        .withNetworkAliases(MYSQL_CONTAINER_HOST)
+                        
.withStartupTimeoutSeconds(MYSQL_STARTUP_TIMEOUT_SECONDS)
+                        .waitingFor(Wait.forHealthcheck())
+                        .withLogConsumer(
+                                new 
Slf4jLogConsumer(DockerLoggerFactory.getLogger(MYSQL_IMAGE)));
+
+        Startables.deepStart(Stream.of(mysqlContainer)).join();
+    }
+
+    @Override
+    @BeforeAll
+    public void startUp() throws Exception {
+        initContainer();
+        initializeJdbcTable();
+    }
+
+    @Override
+    @AfterAll
+    public void tearDown() {
+        if (mysqlContainer != null) {
+            mysqlContainer.close();
+        }
+    }
+
+    @TestTemplate
+    public void testTableOptionsSink(TestContainer container)
+            throws IOException, InterruptedException, SQLException {
+        try {
+            Container.ExecResult execResult = 
container.executeJob(CONFIG_FILE);
+            Assertions.assertEquals(0, execResult.getExitCode(), 
execResult.getStderr());
+            assertSinkTableOptions();
+        } finally {
+            clearSinkTable();
+        }
+    }
+
+    private void assertSinkTableOptions() throws SQLException {
+        try (Connection connection = getJdbcConnection();
+                Statement statement = connection.createStatement()) {
+            ResultSet createTableResult =
+                    statement.executeQuery(
+                            String.format(
+                                    "SHOW CREATE TABLE `%s`.`%s`", 
MYSQL_DATABASE, MYSQL_SINK));
+            Assertions.assertTrue(createTableResult.next());
+            String createTableSql = 
createTableResult.getString(2).toLowerCase();
+            Assertions.assertTrue(
+                    createTableSql.contains(
+                            MySqlCatalog.TABLE_OPTION_ENGINE.toLowerCase() + 
"=innodb"),
+                    createTableSql);
+            Assertions.assertTrue(
+                    createTableSql.contains(
+                            MySqlCatalog.TABLE_OPTION_CHARSET.toLowerCase() + 
"=utf8mb4"),
+                    createTableSql);
+            Assertions.assertTrue(
+                    createTableSql.contains(
+                            MySqlCatalog.TABLE_OPTION_COLLATE.toLowerCase()
+                                    + "=utf8mb4_unicode_ci"),
+                    createTableSql);
+
+            ResultSet countResult =
+                    statement.executeQuery(
+                            String.format(
+                                    "SELECT COUNT(*) FROM `%s`.`%s`", 
MYSQL_DATABASE, MYSQL_SINK));
+            Assertions.assertTrue(countResult.next());
+            Assertions.assertEquals(3, countResult.getInt(1));
+        }
+    }
+
+    private Connection getJdbcConnection() throws SQLException {
+        return DriverManager.getConnection(
+                mysqlContainer.getJdbcUrl(),
+                mysqlContainer.getUsername(),
+                mysqlContainer.getPassword());
+    }
+
+    private void initializeJdbcTable() {
+        try (Connection connection = getJdbcConnection();
+                Statement statement = connection.createStatement()) {
+            statement.execute(CREATE_SOURCE_TABLE_SQL);
+            statement.execute(INSERT_SOURCE_SQL);
+        } catch (SQLException e) {
+            throw new RuntimeException("Initializing MySQL table failed!", e);
+        }
+    }
+
+    private void clearSinkTable() {
+        try (Connection connection = getJdbcConnection();
+                Statement statement = connection.createStatement()) {
+            statement.execute(
+                    String.format("DROP TABLE IF EXISTS `%s`.`%s`", 
MYSQL_DATABASE, MYSQL_SINK));
+        } catch (SQLException e) {
+            throw new RuntimeException("Clearing sink table failed!", e);
+        }
+    }
+}
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-7/src/test/resources/jdbc_mysql_sink_with_table_options.conf
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-7/src/test/resources/jdbc_mysql_sink_with_table_options.conf
new file mode 100644
index 0000000000..99b09f3f6e
--- /dev/null
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-7/src/test/resources/jdbc_mysql_sink_with_table_options.conf
@@ -0,0 +1,54 @@
+#
+# 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.
+#
+
+env {
+  parallelism = 1
+  job.mode = "BATCH"
+}
+
+source {
+  jdbc {
+    url = "jdbc:mysql://mysql-e2e-table-options:3306/seatunnel?useSSL=false"
+    driver = "com.mysql.cj.jdbc.Driver"
+    connection_check_timeout_sec = 100
+    username = "root"
+    password = "Abc!@#135_seatunnel"
+    query = "select * from source;"
+  }
+}
+
+sink {
+  jdbc {
+    url = "jdbc:mysql://mysql-e2e-table-options:3306/seatunnel?useSSL=false"
+    driver = "com.mysql.cj.jdbc.Driver"
+    username = "root"
+    password = "Abc!@#135_seatunnel"
+
+    generate_sink_sql = true
+    database = "seatunnel"
+    table = "sink_table_options"
+    primary_keys = ["id"]
+
+    schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST"
+    data_save_mode = "APPEND_DATA"
+    table_options = {
+      "engine" = "InnoDB"
+      "charset" = "utf8mb4"
+      "collate" = "utf8mb4_unicode_ci"
+    }
+  }
+}


Reply via email to