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");