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 f274e165e4 [Feature][Connector-V2] Support JDBC sink table_options for 
Dameng (#11419)
f274e165e4 is described below

commit f274e165e40781e2b944d4bac69cfa8c46a52a3a
Author: luxiaolong <[email protected]>
AuthorDate: Thu Jul 30 23:50:10 2026 +0800

    [Feature][Connector-V2] Support JDBC sink table_options for Dameng (#11419)
    
    Co-authored-by: det101 <[email protected]>
    Co-authored-by: Cursor <[email protected]>
---
 docs/en/connectors/sink/Jdbc.md                    |  28 +++-
 docs/zh/connectors/sink/Jdbc.md                    |  28 +++-
 .../seatunnel/jdbc/catalog/dm/DamengCatalog.java   |   9 ++
 .../catalog/dm/DamengCreateTableSqlBuilder.java    |  18 +++
 .../jdbc/internal/dialect/dm/DmdbDialect.java      | 127 ++++++++++++++++++
 .../dm/DamengCreateTableSqlBuilderTest.java        | 144 +++++++++++++++++++++
 .../jdbc/internal/dialect/dm/DmdbDialectTest.java  | 115 ++++++++++++++++
 .../JdbcTableOptionsConditionExtensionTest.java    |  33 +++++
 8 files changed, 500 insertions(+), 2 deletions(-)

diff --git a/docs/en/connectors/sink/Jdbc.md b/docs/en/connectors/sink/Jdbc.md
index ab927008e8..3622c604ca 100644
--- a/docs/en/connectors/sink/Jdbc.md
+++ b/docs/en/connectors/sink/Jdbc.md
@@ -382,6 +382,7 @@ Current support:
 | TiDB | Yes | `engine`, `charset`, `collate` (via MySQL JDBC protocol and 
`jdbc:mysql://`) |
 | OceanBase (MySQL mode) | Yes | `engine`, `charset`, `collate` |
 | PostgreSQL | Yes | `tablespace`, `fillfactor` |
+| Dameng | Yes | `tablespace`, `fillfactor` |
 | Oracle | Yes | `tablespace`, `pctfree` |
 | OceanBase (Oracle mode) | Yes | `tablespace`, `pctfree` (via 
`compatible_mode=oracle` → Oracle dialect / DDL path) |
 | Kingbase | Yes | `tablespace`, `fillfactor` |
@@ -395,10 +396,11 @@ Invalid or unsupported keys are validated early via 
`JdbcSinkFactory` option rul
 - **TiDB**: When connected via `jdbc:mysql://` with a MySQL JDBC driver, TiDB 
shares the same key whitelist and DDL merge path as MySQL. `charset` and 
`collate` take effect; `engine` is accepted for MySQL syntax compatibility but 
is **ignored** by TiDB (storage engine is not configurable).
 - **OceanBase (MySQL mode)**: Supported for `jdbc:oceanbase://` when not using 
Oracle-compatible mode. `charset` and `collate` must be values supported by 
your OceanBase version (typically a MySQL-compatible subset; use `SHOW CHARSET` 
/ `SHOW COLLATION` on the target). Unsupported values fail when `CREATE TABLE` 
runs, not at job submission.
 - **PostgreSQL**: `fillfactor` is emitted as `WITH (fillfactor=<n>)` and must 
be an integer in `[10, 100]`; `tablespace` is emitted as `TABLESPACE "..."` 
using the configured name literally (not rewritten by `fieldIde`). Blank values 
and illegal characters in `tablespace` (for example `"`) are rejected at job 
submission. Only these curated keys are accepted (arbitrary `WITH` parameters 
are not supported). OpenGauss and HighGo inherit the same validation and DDL 
path via Postgres catalog/ [...]
+- **Dameng**: `fillfactor` and `tablespace` are emitted as a Dameng `STORAGE 
(...)` clause (`FILLFACTOR <n>`, `ON "<tablespace>"`). `fillfactor` must be an 
integer in `[0, 100]`; `tablespace` uses the configured name literally (not 
rewritten by `fieldIde`). Blank values and illegal characters in `tablespace` 
(for example `"`) are rejected at job submission. Only these curated keys are 
accepted (arbitrary nested `STORAGE` parameters such as `INITIAL` / `NEXT` are 
not supported). The table [...]
 - **Oracle / OceanBase (Oracle mode)**: `pctfree` is emitted as `PCTFREE <n>` 
and must be an integer in `[0, 99]`; `tablespace` is emitted as `TABLESPACE 
"..."` using the configured name literally (not rewritten by `fieldIde`). Blank 
values and illegal characters in `tablespace` (for example `"`) are rejected at 
job submission. Only these curated keys are accepted (nested `STORAGE (...)` 
and LOB/partition clauses are not supported). The tablespace must already exist 
on the target.
 - **Kingbase**: `fillfactor` is emitted as `WITH (fillfactor=<n>)` and must be 
an integer in `[10, 100]` (PostgreSQL-compatible); `tablespace` is emitted as 
`TABLESPACE "..."` using the configured name literally (not rewritten by 
`fieldIde`). Blank values and illegal characters in `tablespace` (for example 
`"`) are rejected at job submission. Only these curated keys are accepted 
(arbitrary `WITH (...)` parameters are not supported). The tablespace must 
already exist on the target.
 
-SeaTunnel validates the **key whitelist** at submission time for all dialects 
that support `table_options`. For PostgreSQL (and OpenGauss / HighGo via the 
same path), and Kingbase, it also validates blank values and the `fillfactor` 
numeric range. For Oracle / OceanBase (Oracle mode), it also validates blank 
values and the `pctfree` numeric range. Other dialects (for example MySQL) do 
not verify whether each value is supported by the target database beyond the 
key whitelist.
+SeaTunnel validates the **key whitelist** at submission time for all dialects 
that support `table_options`. For PostgreSQL (and OpenGauss / HighGo via the 
same path), Dameng, and Kingbase, it also validates blank values and the 
`fillfactor` numeric range. For Oracle / OceanBase (Oracle mode), it also 
validates blank values and the `pctfree` numeric range. Other dialects (for 
example MySQL) do not verify whether each value is supported by the target 
database beyond the key whitelist.
 
 Example (MySQL auto-create with engine and charset):
 
@@ -447,6 +449,30 @@ sink {
 }
 ```
 
+Example (Dameng auto-create with tablespace and fillfactor):
+
+```hocon
+sink {
+  Jdbc {
+    url = "jdbc:dm://localhost:5236"
+    driver = "dm.jdbc.driver.DmDriver"
+    username = "SYSDBA"
+    password = "SYSDBA"
+    database = "DAMENG"
+    table = "orders"
+    generate_sink_sql = true
+    schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST"
+    primary_keys = ["id"]
+    table_options = {
+      "tablespace" = "MAIN"
+      "fillfactor" = "80"
+    }
+  }
+}
+```
+
+The generated `CREATE TABLE` statement appends `STORAGE (FILLFACTOR 80, ON 
"MAIN")`.
+
 Example (Oracle auto-create with tablespace and pctfree):
 
 ```hocon
diff --git a/docs/zh/connectors/sink/Jdbc.md b/docs/zh/connectors/sink/Jdbc.md
index 01b28eb4e4..c486d61336 100644
--- a/docs/zh/connectors/sink/Jdbc.md
+++ b/docs/zh/connectors/sink/Jdbc.md
@@ -373,6 +373,7 @@ Sink 在自动建表(SaveMode DDL)时附加的表级选项。仅在 `schema_
 | TiDB | 是 | `engine`、`charset`、`collate`(通过 MySQL JDBC 协议与 `jdbc:mysql://` 
连接) |
 | OceanBase(MySQL 模式) | 是 | `engine`、`charset`、`collate` |
 | PostgreSQL | 是 | `tablespace`、`fillfactor` |
+| 达梦 Dameng | 是 | `tablespace`、`fillfactor` |
 | Oracle | 是 | `tablespace`、`pctfree` |
 | OceanBase(Oracle 模式) | 是 | `tablespace`、`pctfree`(`compatible_mode=oracle` 
时走 Oracle 方言 / DDL 路径) |
 | Kingbase | 是 | `tablespace`、`fillfactor` |
@@ -386,10 +387,11 @@ Sink 在自动建表(SaveMode DDL)时附加的表级选项。仅在 `schema_
 - **TiDB**:通过 `jdbc:mysql://` 与 MySQL JDBC 驱动连接时,与 MySQL 使用相同的 key 白名单与 DDL 
拼接方式。`charset`、`collate` 会生效;`engine` 仅为 MySQL 兼容语法,TiDB 会解析但**忽略**存储引擎设置。
 - **OceanBase(MySQL 模式)**:`jdbc:oceanbase://` 且非 Oracle 兼容模式时支持上述三个 
key。`charset`、`collate` 须为 OceanBase 当前版本支持的字符集与排序规则(通常为 MySQL 兼容子集,请以目标库 `SHOW 
CHARSET` / `SHOW COLLATION` 为准);不支持的取值会在执行 `CREATE TABLE` 时报错,而非在作业提交阶段校验。
 - **PostgreSQL**:`fillfactor` 会生成 `WITH (fillfactor=<n>)`,取值须为 `[10, 100]` 
的整数;`tablespace` 会生成 `TABLESPACE "..."`,按配置字面量引用(**不受** `fieldIde` 
大小写改写)。空白值,以及 `tablespace` 中的非法字符(例如 `"`)会在作业提交阶段被拒绝。仅接受上述 curated key(不支持任意 
`WITH` 参数)。OpenGauss、HighGo 通过 Postgres catalog/dialect 继承同一套校验与 DDL。
+- **达梦 Dameng**:`fillfactor` 与 `tablespace` 会写入达梦 `STORAGE (...)` 
子句(`FILLFACTOR <n>`、`ON "<表空间>"`)。`fillfactor` 取值须为 `[0, 100]` 的整数;`tablespace` 
按配置字面量引用(**不受** `fieldIde` 大小写改写)。空白值,以及 `tablespace` 中的非法字符(例如 
`"`)会在作业提交阶段被拒绝。仅接受上述 curated key(不支持透传任意 `STORAGE` 参数如 `INITIAL` / 
`NEXT`)。表空间须已在目标库存在。
 - **Oracle / OceanBase(Oracle 模式)**:`pctfree` 会生成 `PCTFREE <n>`,取值须为 `[0, 99]` 
的整数;`tablespace` 会生成 `TABLESPACE "..."`,按配置字面量引用(**不受** `fieldIde` 
大小写改写)。空白值,以及 `tablespace` 中的非法字符(例如 `"`)会在作业提交阶段被拒绝。仅接受上述 curated key(不支持嵌套 
`STORAGE (...)`、LOB/分区子句)。目标库中 tablespace 须事先存在。
 - **Kingbase**:`fillfactor` 会写入 `WITH (fillfactor=<n>)`,取值须为 `[10, 100]` 
的整数(PostgreSQL 兼容);`tablespace` 会写入 `TABLESPACE "..."`,按配置字面量引用(**不受** 
`fieldIde` 大小写改写)。空白值,以及 `tablespace` 中的非法字符(例如 `"`)会在作业提交阶段被拒绝。仅接受上述 curated 
key(不支持透传任意 `WITH (...)` 参数)。表空间须已在目标库存在。
 
-SeaTunnel 在提交时会对所有支持 `table_options` 的方言校验 **key 白名单**。对 PostgreSQL(以及 
OpenGauss / HighGo 同源路径)与 Kingbase,还会校验空白值与 `fillfactor` 数值区间。对 Oracle / 
OceanBase(Oracle 模式),还会校验空白值与 `pctfree` 数值区间。其他方言(例如 
MySQL)在白名单之外不额外校验具体取值是否被目标库支持。
+SeaTunnel 在提交时会对所有支持 `table_options` 的方言校验 **key 白名单**。对 PostgreSQL(以及 
OpenGauss / HighGo 同源路径)、达梦与 Kingbase,还会校验空白值与 `fillfactor` 数值区间。对 Oracle / 
OceanBase(Oracle 模式),还会校验空白值与 `pctfree` 数值区间。其他方言(例如 
MySQL)在白名单之外不额外校验具体取值是否被目标库支持。
 
 示例(MySQL 自动建表时指定存储引擎与字符集):
 
@@ -438,6 +440,30 @@ sink {
 }
 ```
 
+示例(达梦自动建表时指定表空间与填充比例):
+
+```hocon
+sink {
+  Jdbc {
+    url = "jdbc:dm://localhost:5236"
+    driver = "dm.jdbc.driver.DmDriver"
+    username = "SYSDBA"
+    password = "SYSDBA"
+    database = "DAMENG"
+    table = "orders"
+    generate_sink_sql = true
+    schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST"
+    primary_keys = ["id"]
+    table_options = {
+      "tablespace" = "MAIN"
+      "fillfactor" = "80"
+    }
+  }
+}
+```
+
+生成的 DDL 会追加 `STORAGE (FILLFACTOR 80, ON "MAIN")`。
+
 示例(Oracle 自动建表时指定 tablespace 与 pctfree):
 
 ```hocon
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/dm/DamengCatalog.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/dm/DamengCatalog.java
index 55dc298165..69dc9f154f 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/dm/DamengCatalog.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/dm/DamengCatalog.java
@@ -42,6 +42,15 @@ import java.util.List;
 @Slf4j
 public class DamengCatalog extends AbstractJdbcCatalog {
 
+    /** Sink {@code table_options} key for Dameng table tablespace ({@code 
STORAGE (ON ...)}). */
+    public static final String TABLE_OPTION_TABLESPACE = "tablespace";
+
+    /**
+     * Sink {@code table_options} key for Dameng page fill factor ({@code 
STORAGE (FILLFACTOR
+     * ...)}).
+     */
+    public static final String TABLE_OPTION_FILLFACTOR = "fillfactor";
+
     private static final String SELECT_COLUMNS_SQL =
             "SELECT COLUMNS.COLUMN_NAME, COLUMNS.DATA_TYPE, 
COLUMNS.DATA_LENGTH, COLUMNS.DATA_PRECISION, COLUMNS.DATA_SCALE "
                     + ", COLUMNS.NULLABLE, COLUMNS.DATA_DEFAULT, 
COMMENTS.COMMENTS ,"
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/dm/DamengCreateTableSqlBuilder.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/dm/DamengCreateTableSqlBuilder.java
index 38b406f785..a8b9e72c1d 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/dm/DamengCreateTableSqlBuilder.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/dm/DamengCreateTableSqlBuilder.java
@@ -27,6 +27,7 @@ import org.apache.seatunnel.api.table.catalog.TablePath;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.catalog.AbstractJdbcCreateTableSqlBuilder;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.catalog.utils.CatalogUtils;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.DatabaseIdentifier;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.dm.DmdbDialect;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.dm.DmdbTypeConverter;
 
 import org.apache.commons.collections4.CollectionUtils;
@@ -41,6 +42,8 @@ public class DamengCreateTableSqlBuilder extends 
AbstractJdbcCreateTableSqlBuild
     private final PrimaryKey primaryKey;
     private final String sourceCatalogName;
     private final String fieldIde;
+    private final String tablespace;
+    private final String fillfactor;
     private final List<ConstraintKey> constraintKeys;
     private boolean createIndex;
 
@@ -49,6 +52,8 @@ public class DamengCreateTableSqlBuilder extends 
AbstractJdbcCreateTableSqlBuild
         this.primaryKey = catalogTable.getTableSchema().getPrimaryKey();
         this.sourceCatalogName = catalogTable.getCatalogName();
         this.fieldIde = catalogTable.getOptions().get("fieldIde");
+        this.tablespace = 
catalogTable.getOptions().get(DamengCatalog.TABLE_OPTION_TABLESPACE);
+        this.fillfactor = 
catalogTable.getOptions().get(DamengCatalog.TABLE_OPTION_FILLFACTOR);
         constraintKeys = catalogTable.getTableSchema().getConstraintKeys();
         this.createIndex = createIndex;
     }
@@ -93,6 +98,19 @@ public class DamengCreateTableSqlBuilder extends 
AbstractJdbcCreateTableSqlBuild
 
         createTableSql.append(String.join(",\n", columnSqls));
         createTableSql.append("\n)");
+        List<String> storageItems = new ArrayList<>();
+        if (StringUtils.isNotBlank(fillfactor)) {
+            storageItems.add("FILLFACTOR " + 
DmdbDialect.normalizeFillfactorForDdl(fillfactor));
+        }
+        if (StringUtils.isNotBlank(tablespace)) {
+            storageItems.add("ON \"" + 
DmdbDialect.normalizeTablespaceForDdl(tablespace) + "\"");
+        }
+        if (!storageItems.isEmpty()) {
+            createTableSql
+                    .append("\nSTORAGE (")
+                    .append(String.join(", ", storageItems))
+                    .append(")");
+        }
         sqls.add(createTableSql.toString());
 
         List<String> commentSqls =
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 a4ca78b7b7..cf9b29baaa 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
@@ -19,6 +19,7 @@ package 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.dm;
 
 import org.apache.seatunnel.shade.org.apache.commons.lang3.StringUtils;
 
+import org.apache.seatunnel.api.common.SeaTunnelAPIErrorCode;
 import org.apache.seatunnel.api.table.catalog.Column;
 import org.apache.seatunnel.api.table.catalog.TablePath;
 import org.apache.seatunnel.api.table.converter.BasicTypeDefine;
@@ -26,6 +27,8 @@ import org.apache.seatunnel.api.table.converter.TypeConverter;
 import org.apache.seatunnel.api.table.schema.event.AlterTableAddColumnEvent;
 import org.apache.seatunnel.api.table.schema.event.AlterTableChangeColumnEvent;
 import org.apache.seatunnel.api.table.schema.event.AlterTableModifyColumnEvent;
+import org.apache.seatunnel.connectors.seatunnel.jdbc.catalog.dm.DamengCatalog;
+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;
@@ -39,8 +42,12 @@ import java.sql.SQLException;
 import java.sql.Statement;
 import java.util.ArrayList;
 import java.util.Arrays;
+import java.util.Collections;
+import java.util.LinkedHashSet;
 import java.util.List;
+import java.util.Map;
 import java.util.Optional;
+import java.util.Set;
 import java.util.stream.Collectors;
 
 import static 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.dm.DmdbTypeConverter.DM_CHAR;
@@ -56,6 +63,18 @@ import static 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.dm
 @Slf4j
 public class DmdbDialect implements JdbcDialect {
 
+    /** Dameng FILLFACTOR legal range (inclusive). */
+    private static final int FILLFACTOR_MIN = 0;
+
+    private static final int FILLFACTOR_MAX = 100;
+
+    private static final Set<String> SUPPORTED_TABLE_OPTIONS =
+            Collections.unmodifiableSet(
+                    new LinkedHashSet<>(
+                            Arrays.asList(
+                                    DamengCatalog.TABLE_OPTION_TABLESPACE,
+                                    DamengCatalog.TABLE_OPTION_FILLFACTOR)));
+
     public String fieldIde;
 
     public DmdbDialect(String fieldIde) {
@@ -386,4 +405,112 @@ public class DmdbDialect implements JdbcDialect {
             return rs.getString("NULLABLE").equals("Y");
         }
     }
+
+    @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)));
+        }
+
+        for (Map.Entry<String, String> entry : tableOptions.entrySet()) {
+            String key = entry.getKey();
+            String value = entry.getValue();
+            if (StringUtils.isBlank(value)) {
+                throw new JdbcConnectorException(
+                        SeaTunnelAPIErrorCode.CONFIG_VALIDATION_FAILED,
+                        String.format(
+                                "Invalid JDBC table_options for dialect '%s': 
key '%s' must not be blank",
+                                dialectName(), key));
+            }
+            String trimmed = value.trim();
+            if (DamengCatalog.TABLE_OPTION_FILLFACTOR.equals(key)) {
+                normalizeFillfactorForDdl(trimmed);
+            } else if (DamengCatalog.TABLE_OPTION_TABLESPACE.equals(key)) {
+                normalizeTablespaceForDdl(trimmed);
+            }
+        }
+    }
+
+    /**
+     * Normalizes a Dameng {@code fillfactor} option for DDL emission.
+     *
+     * <p>Shared by submission-time validation and {@link
+     * 
org.apache.seatunnel.connectors.seatunnel.jdbc.catalog.dm.DamengCreateTableSqlBuilder}
 so
+     * catalog-driven auto-create cannot bypass the value contract.
+     */
+    public static String normalizeFillfactorForDdl(String value) {
+        return Integer.toString(parseFillfactorInRange(value));
+    }
+
+    /**
+     * Normalizes a Dameng {@code tablespace} option for DDL emission.
+     *
+     * <p>Shared by submission-time validation and {@link
+     * 
org.apache.seatunnel.connectors.seatunnel.jdbc.catalog.dm.DamengCreateTableSqlBuilder}.
+     */
+    public static String normalizeTablespaceForDdl(String value) {
+        String trimmed = value.trim();
+        validateTablespaceCharacters(trimmed);
+        return trimmed;
+    }
+
+    private static int parseFillfactorInRange(String value) {
+        int fillfactor;
+        try {
+            fillfactor = Integer.parseInt(value);
+        } catch (NumberFormatException e) {
+            throw new JdbcConnectorException(
+                    SeaTunnelAPIErrorCode.CONFIG_VALIDATION_FAILED,
+                    String.format(
+                            "Invalid JDBC table_options for dialect '%s': key 
'%s' must be an integer between %d and %d, but got '%s'",
+                            DatabaseIdentifier.DAMENG,
+                            DamengCatalog.TABLE_OPTION_FILLFACTOR,
+                            FILLFACTOR_MIN,
+                            FILLFACTOR_MAX,
+                            value));
+        }
+        if (fillfactor < FILLFACTOR_MIN || fillfactor > FILLFACTOR_MAX) {
+            throw new JdbcConnectorException(
+                    SeaTunnelAPIErrorCode.CONFIG_VALIDATION_FAILED,
+                    String.format(
+                            "Invalid JDBC table_options for dialect '%s': key 
'%s' must be an integer between %d and %d, but got '%s'",
+                            DatabaseIdentifier.DAMENG,
+                            DamengCatalog.TABLE_OPTION_FILLFACTOR,
+                            FILLFACTOR_MIN,
+                            FILLFACTOR_MAX,
+                            value));
+        }
+        return fillfactor;
+    }
+
+    private static void validateTablespaceCharacters(String value) {
+        if (value.indexOf('"') >= 0 || value.indexOf(';') >= 0) {
+            throw illegalTablespaceException(value);
+        }
+        for (int i = 0; i < value.length(); i++) {
+            if (Character.isISOControl(value.charAt(i))) {
+                throw illegalTablespaceException(value);
+            }
+        }
+    }
+
+    private static JdbcConnectorException illegalTablespaceException(String 
value) {
+        return new JdbcConnectorException(
+                SeaTunnelAPIErrorCode.CONFIG_VALIDATION_FAILED,
+                String.format(
+                        "Invalid JDBC table_options for dialect '%s': key '%s' 
contains illegal characters: '%s'",
+                        DatabaseIdentifier.DAMENG, 
DamengCatalog.TABLE_OPTION_TABLESPACE, value));
+    }
 }
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/dm/DamengCreateTableSqlBuilderTest.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/dm/DamengCreateTableSqlBuilderTest.java
index a6718cf00a..14c0632b5c 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/dm/DamengCreateTableSqlBuilderTest.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/dm/DamengCreateTableSqlBuilderTest.java
@@ -28,6 +28,7 @@ import org.apache.seatunnel.api.table.catalog.TablePath;
 import org.apache.seatunnel.api.table.catalog.TableSchema;
 import org.apache.seatunnel.api.table.type.BasicType;
 import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.exception.JdbcConnectorException;
 
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
@@ -193,6 +194,149 @@ public class DamengCreateTableSqlBuilderTest {
         Assertions.assertTrue(sqls.get(0).startsWith("CREATE TABLE"));
     }
 
+    @Test
+    public void testBuildCreateTableSqlWithTableOptions() {
+        TablePath tablePath = TablePath.of("test_database", "test_schema", 
"test_table");
+        TableSchema tableSchema =
+                TableSchema.builder()
+                        .column(PhysicalColumn.of("id", BasicType.LONG_TYPE, 
22, false, null, null))
+                        .build();
+
+        HashMap<String, String> options = new HashMap<>();
+        options.put(DamengCatalog.TABLE_OPTION_TABLESPACE, "MAIN");
+        options.put(DamengCatalog.TABLE_OPTION_FILLFACTOR, "80");
+
+        CatalogTable catalogTable =
+                CatalogTable.of(
+                        TableIdentifier.of("test_catalog", tablePath),
+                        tableSchema,
+                        options,
+                        new ArrayList<>(),
+                        "Table with storage options");
+
+        List<String> sqls = new DamengCreateTableSqlBuilder(catalogTable, 
false).build(tablePath);
+
+        Assertions.assertEquals(1, sqls.size());
+        String createTable = sqls.get(0);
+        Assertions.assertTrue(
+                createTable.contains("STORAGE (FILLFACTOR 80, ON \"MAIN\")"),
+                "CREATE TABLE should contain Dameng STORAGE clause: " + 
createTable);
+    }
+
+    @Test
+    public void testBuildCreateTableSqlWithTableOptionsIgnoresFieldIde() {
+        TablePath tablePath = TablePath.of("test_database", "test_schema", 
"test_table");
+        TableSchema tableSchema =
+                TableSchema.builder()
+                        .column(PhysicalColumn.of("id", BasicType.LONG_TYPE, 
22, false, null, null))
+                        .build();
+
+        HashMap<String, String> options = new HashMap<>();
+        options.put("fieldIde", "LOWERCASE");
+        options.put(DamengCatalog.TABLE_OPTION_TABLESPACE, "MAIN");
+        options.put(DamengCatalog.TABLE_OPTION_FILLFACTOR, "80");
+
+        CatalogTable catalogTable =
+                CatalogTable.of(
+                        TableIdentifier.of("test_catalog", tablePath),
+                        tableSchema,
+                        options,
+                        new ArrayList<>(),
+                        "Table with storage options");
+
+        List<String> sqls = new DamengCreateTableSqlBuilder(catalogTable, 
false).build(tablePath);
+
+        Assertions.assertEquals(1, sqls.size());
+        String createTable = sqls.get(0);
+        Assertions.assertTrue(
+                createTable.contains("STORAGE (FILLFACTOR 80, ON \"MAIN\")"),
+                "tablespace must not be rewritten by fieldIde; got: " + 
createTable);
+        Assertions.assertFalse(createTable.contains("ON \"main\""));
+    }
+
+    @Test
+    public void testBuildCreateTableSqlNormalizesFillfactorInDdl() {
+        TablePath tablePath = TablePath.of("test_database", "test_schema", 
"test_table");
+        TableSchema tableSchema =
+                TableSchema.builder()
+                        .column(PhysicalColumn.of("id", BasicType.LONG_TYPE, 
22, false, null, null))
+                        .build();
+
+        HashMap<String, String> options = new HashMap<>();
+        options.put(DamengCatalog.TABLE_OPTION_FILLFACTOR, "+80");
+
+        CatalogTable catalogTable =
+                CatalogTable.of(
+                        TableIdentifier.of("test_catalog", tablePath),
+                        tableSchema,
+                        options,
+                        new ArrayList<>(),
+                        "Table with normalized fillfactor");
+
+        List<String> sqls = new DamengCreateTableSqlBuilder(catalogTable, 
false).build(tablePath);
+
+        Assertions.assertEquals(1, sqls.size());
+        Assertions.assertTrue(
+                sqls.get(0).contains("STORAGE (FILLFACTOR 80)"),
+                "fillfactor should be normalized in DDL: " + sqls.get(0));
+    }
+
+    @Test
+    public void 
testBuildCreateTableSqlRejectsInvalidFillfactorWithoutSubmissionValidation() {
+        TablePath tablePath = TablePath.of("test_database", "test_schema", 
"test_table");
+        TableSchema tableSchema =
+                TableSchema.builder()
+                        .column(PhysicalColumn.of("id", BasicType.LONG_TYPE, 
22, false, null, null))
+                        .build();
+
+        HashMap<String, String> options = new HashMap<>();
+        options.put(DamengCatalog.TABLE_OPTION_FILLFACTOR, "abc");
+
+        CatalogTable catalogTable =
+                CatalogTable.of(
+                        TableIdentifier.of("test_catalog", tablePath),
+                        tableSchema,
+                        options,
+                        new ArrayList<>(),
+                        "Table with invalid fillfactor");
+
+        JdbcConnectorException exception =
+                Assertions.assertThrows(
+                        JdbcConnectorException.class,
+                        () ->
+                                new DamengCreateTableSqlBuilder(catalogTable, 
false)
+                                        .build(tablePath));
+        Assertions.assertTrue(exception.getMessage().contains("must be an 
integer between"));
+    }
+
+    @Test
+    public void 
testBuildCreateTableSqlRejectsIllegalTablespaceWithoutSubmissionValidation() {
+        TablePath tablePath = TablePath.of("test_database", "test_schema", 
"test_table");
+        TableSchema tableSchema =
+                TableSchema.builder()
+                        .column(PhysicalColumn.of("id", BasicType.LONG_TYPE, 
22, false, null, null))
+                        .build();
+
+        HashMap<String, String> options = new HashMap<>();
+        options.put(DamengCatalog.TABLE_OPTION_TABLESPACE, "MAIN\"TS");
+
+        CatalogTable catalogTable =
+                CatalogTable.of(
+                        TableIdentifier.of("test_catalog", tablePath),
+                        tableSchema,
+                        options,
+                        new ArrayList<>(),
+                        "Table with illegal tablespace");
+
+        JdbcConnectorException exception =
+                Assertions.assertThrows(
+                        JdbcConnectorException.class,
+                        () ->
+                                new DamengCreateTableSqlBuilder(catalogTable, 
false)
+                                        .build(tablePath));
+        Assertions.assertTrue(exception.getMessage().contains("illegal 
characters"));
+    }
+
     @Test
     public void testColumnSinkType() {
         DamengCreateTableSqlBuilder sqlBuilder = 
mock(DamengCreateTableSqlBuilder.class);
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 bc7a81f14d..b1cacdc51b 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
@@ -17,12 +17,17 @@
 
 package org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.dm;
 
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.exception.JdbcConnectorException;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.JdbcDialect;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.dialectenum.FieldIdeEnum;
 
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
+
 public class DmdbDialectTest {
     @Test
     public void testIdentifierCaseSensitive() {
@@ -44,4 +49,114 @@ public class DmdbDialectTest {
         Assertions.assertEquals("\"TEST\"", dialect.quoteIdentifier("test"));
         Assertions.assertEquals("\"TEST\"", dialect.quoteIdentifier("TEST"));
     }
+
+    @Test
+    void testValidateTableOptionsAcceptsSupportedKeys() {
+        DmdbDialect dialect = new 
DmdbDialect(FieldIdeEnum.ORIGINAL.getValue());
+        Map<String, String> tableOptions = new HashMap<>();
+        tableOptions.put("tablespace", "MAIN");
+        tableOptions.put("fillfactor", "80");
+        Assertions.assertDoesNotThrow(() -> 
dialect.validateTableOptions(tableOptions));
+    }
+
+    @Test
+    void testValidateTableOptionsFillfactorBoundary() {
+        DmdbDialect dialect = new 
DmdbDialect(FieldIdeEnum.ORIGINAL.getValue());
+        Assertions.assertDoesNotThrow(
+                () -> 
dialect.validateTableOptions(Collections.singletonMap("fillfactor", "0")));
+        Assertions.assertDoesNotThrow(
+                () -> 
dialect.validateTableOptions(Collections.singletonMap("fillfactor", "100")));
+    }
+
+    @Test
+    void testValidateTableOptionsRejectsUnsupportedKeys() {
+        DmdbDialect dialect = new 
DmdbDialect(FieldIdeEnum.ORIGINAL.getValue());
+        Map<String, String> tableOptions = new HashMap<>();
+        tableOptions.put("pctfree", "10");
+        JdbcConnectorException exception =
+                Assertions.assertThrows(
+                        JdbcConnectorException.class,
+                        () -> dialect.validateTableOptions(tableOptions));
+        Assertions.assertTrue(exception.getMessage().contains("Unsupported 
JDBC table_options"));
+        Assertions.assertTrue(exception.getMessage().contains("Dameng"));
+    }
+
+    @Test
+    void testValidateTableOptionsRejectBlankValues() {
+        DmdbDialect dialect = new 
DmdbDialect(FieldIdeEnum.ORIGINAL.getValue());
+
+        JdbcConnectorException blankTablespace =
+                Assertions.assertThrows(
+                        JdbcConnectorException.class,
+                        () ->
+                                dialect.validateTableOptions(
+                                        Collections.singletonMap("tablespace", 
" ")));
+        Assertions.assertTrue(blankTablespace.getMessage().contains("must not 
be blank"));
+
+        JdbcConnectorException blankFillfactor =
+                Assertions.assertThrows(
+                        JdbcConnectorException.class,
+                        () ->
+                                dialect.validateTableOptions(
+                                        Collections.singletonMap("fillfactor", 
" ")));
+        Assertions.assertTrue(blankFillfactor.getMessage().contains("must not 
be blank"));
+    }
+
+    @Test
+    void testValidateTableOptionsRejectInvalidFillfactor() {
+        DmdbDialect dialect = new 
DmdbDialect(FieldIdeEnum.ORIGINAL.getValue());
+
+        JdbcConnectorException nonNumeric =
+                Assertions.assertThrows(
+                        JdbcConnectorException.class,
+                        () ->
+                                dialect.validateTableOptions(
+                                        Collections.singletonMap("fillfactor", 
"abc")));
+        Assertions.assertTrue(nonNumeric.getMessage().contains("must be an 
integer between"));
+
+        JdbcConnectorException outOfRange =
+                Assertions.assertThrows(
+                        JdbcConnectorException.class,
+                        () ->
+                                dialect.validateTableOptions(
+                                        Collections.singletonMap("fillfactor", 
"101")));
+        Assertions.assertTrue(outOfRange.getMessage().contains("must be an 
integer between"));
+    }
+
+    @Test
+    void testValidateTableOptionsRejectIllegalTablespace() {
+        DmdbDialect dialect = new 
DmdbDialect(FieldIdeEnum.ORIGINAL.getValue());
+
+        JdbcConnectorException exception =
+                Assertions.assertThrows(
+                        JdbcConnectorException.class,
+                        () ->
+                                dialect.validateTableOptions(
+                                        Collections.singletonMap("tablespace", 
"MAIN\"TS")));
+        Assertions.assertTrue(exception.getMessage().contains("illegal 
characters"));
+    }
+
+    @Test
+    void testNormalizeFillfactorAcceptsPlusSignAndLeadingZeros() {
+        Assertions.assertEquals("80", 
DmdbDialect.normalizeFillfactorForDdl("+80"));
+        Assertions.assertEquals("80", 
DmdbDialect.normalizeFillfactorForDdl("080"));
+    }
+
+    @Test
+    void testNormalizeFillfactorRejectsOutOfRange() {
+        JdbcConnectorException exception =
+                Assertions.assertThrows(
+                        JdbcConnectorException.class,
+                        () -> DmdbDialect.normalizeFillfactorForDdl("101"));
+        Assertions.assertTrue(exception.getMessage().contains("must be an 
integer between"));
+    }
+
+    @Test
+    void testNormalizeTablespaceRejectsControlCharacters() {
+        JdbcConnectorException tabCharacter =
+                Assertions.assertThrows(
+                        JdbcConnectorException.class,
+                        () -> 
DmdbDialect.normalizeTablespaceForDdl("MAIN\tTS"));
+        Assertions.assertTrue(tabCharacter.getMessage().contains("illegal 
characters"));
+    }
 }
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
index ce7ae851d4..f59e9d2ace 100644
--- 
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
@@ -159,6 +159,31 @@ class JdbcTableOptionsConditionExtensionTest {
         Assertions.assertTrue(exception.getMessage().contains("Oracle"));
     }
 
+    @Test
+    void testDamengTableOptionsPassViaOptionRule() {
+        Map<String, Object> config = damengSinkConfig();
+        Map<String, String> tableOptions = new HashMap<>();
+        tableOptions.put("tablespace", "MAIN");
+        tableOptions.put("fillfactor", "80");
+        config.put(SinkConnectorCommonOptions.TABLE_OPTIONS.key(), 
tableOptions);
+
+        Assertions.assertDoesNotThrow(() -> validateSinkOptionRule(config));
+    }
+
+    @Test
+    void testDamengRejectsUnknownTableOptionsViaOptionRule() {
+        Map<String, Object> config = damengSinkConfig();
+        Map<String, String> tableOptions = new HashMap<>();
+        tableOptions.put("engine", "InnoDB");
+        config.put(SinkConnectorCommonOptions.TABLE_OPTIONS.key(), 
tableOptions);
+
+        OptionValidationException exception =
+                Assertions.assertThrows(
+                        OptionValidationException.class, () -> 
validateSinkOptionRule(config));
+        Assertions.assertTrue(exception.getMessage().contains("Unsupported 
JDBC table_options"));
+        Assertions.assertTrue(exception.getMessage().contains("Dameng"));
+    }
+
     @Test
     void testKingbaseTableOptionsPassViaOptionRule() {
         Map<String, Object> config = kingbaseSinkConfig();
@@ -238,6 +263,14 @@ class JdbcTableOptionsConditionExtensionTest {
         return config;
     }
 
+    private static Map<String, Object> damengSinkConfig() {
+        Map<String, Object> config = new HashMap<>();
+        config.put("url", "jdbc:dm://127.0.0.1:5236");
+        config.put("driver", "dm.jdbc.driver.DmDriver");
+        config.put("query", "INSERT INTO test_table VALUES (?)");
+        return config;
+    }
+
     private static Map<String, Object> kingbaseSinkConfig() {
         Map<String, Object> config = new HashMap<>();
         config.put("url", "jdbc:kingbase8://127.0.0.1:54321/test");

Reply via email to