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 ad7920696f [Feature][Connector-V2] Support JDBC sink table_options for
PostgreSQL (#11417)
ad7920696f is described below
commit ad7920696fa23ea2a87ae6f4700bafc97d1fa16a
Author: luxiaolong <[email protected]>
AuthorDate: Wed Jul 15 22:49:45 2026 +0800
[Feature][Connector-V2] Support JDBC sink table_options for PostgreSQL
(#11417)
Co-authored-by: det101 <[email protected]>
Co-authored-by: Cursor <[email protected]>
---
docs/en/connectors/sink/Jdbc.md | 26 ++++-
docs/zh/connectors/sink/Jdbc.md | 26 ++++-
.../jdbc/catalog/psql/PostgresCatalog.java | 3 +
.../psql/PostgresCreateTableSqlBuilder.java | 16 ++-
.../internal/dialect/psql/PostgresDialect.java | 96 ++++++++++++++++
.../psql/PostgresCreateTableSqlBuilderTest.java | 50 ++++++++
.../psql/PostgresDialectTableOptionsTest.java | 126 +++++++++++++++++++++
.../JdbcTableOptionsConditionExtensionTest.java | 17 ++-
8 files changed, 354 insertions(+), 6 deletions(-)
diff --git a/docs/en/connectors/sink/Jdbc.md b/docs/en/connectors/sink/Jdbc.md
index 0b3b270f98..efd007b6f5 100644
--- a/docs/en/connectors/sink/Jdbc.md
+++ b/docs/en/connectors/sink/Jdbc.md
@@ -270,6 +270,7 @@ Current support:
| MySQL | Yes | `engine`, `charset`, `collate` |
| TiDB | Yes | `engine`, `charset`, `collate` (via MySQL JDBC protocol and
`jdbc:mysql://`) |
| OceanBase (MySQL mode) | Yes | `engine`, `charset`, `collate` |
+| PostgreSQL | Yes | `tablespace`, `fillfactor` |
| 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.
@@ -279,8 +280,9 @@ Invalid or unsupported keys are validated early via
`JdbcSinkFactory` option rul
- **MySQL**: `engine`, `charset`, and `collate` are appended to `CREATE TABLE`
and take effect.
- **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. OceanBase **Oracle-compatible mode** does not
support `table_options`; a non-empty map fails 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/ [...]
-SeaTunnel validates the **key whitelist** at submission time only; it does not
verify whether each value is supported by the target database (same as the
initial MySQL delivery).
+SeaTunnel validates the **key whitelist** at submission time for all dialects
that support `table_options`. For PostgreSQL (and OpenGauss / HighGo via the
same path), it also validates blank values and the `fillfactor` 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):
@@ -307,6 +309,28 @@ sink {
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.
+Example (PostgreSQL auto-create with tablespace and fillfactor):
+
+```hocon
+sink {
+ Jdbc {
+ url = "jdbc:postgresql://localhost:5432/mydb"
+ driver = "org.postgresql.Driver"
+ username = "postgres"
+ password = "password"
+ database = "mydb"
+ table = "public.orders"
+ generate_sink_sql = true
+ schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST"
+ primary_keys = ["id"]
+ table_options = {
+ "tablespace" = "pg_default"
+ "fillfactor" = "70"
+ }
+ }
+}
+```
+
### 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 0df4d5b871..d501df35ff 100644
--- a/docs/zh/connectors/sink/Jdbc.md
+++ b/docs/zh/connectors/sink/Jdbc.md
@@ -258,6 +258,7 @@ Sink 在自动建表(SaveMode DDL)时附加的表级选项。仅在 `schema_
| MySQL | 是 | `engine`、`charset`、`collate` |
| TiDB | 是 | `engine`、`charset`、`collate`(通过 MySQL JDBC 协议与 `jdbc:mysql://`
连接) |
| OceanBase(MySQL 模式) | 是 | `engine`、`charset`、`collate` |
+| PostgreSQL | 是 | `tablespace`、`fillfactor` |
| 其他 JDBC 方言 | 否 | 配置非空 `table_options` 时任务启动即校验失败 |
非法或不支持的 key 会在 `JdbcSinkFactory` 的 option 规则阶段提前校验(`--check` 与作业提交),而非仅在运行时
DDL 阶段失败。
@@ -267,8 +268,9 @@ Sink 在自动建表(SaveMode DDL)时附加的表级选项。仅在 `schema_
- **MySQL**:`engine`、`charset`、`collate` 均会写入 `CREATE TABLE` 并生效。
- **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`
时报错,而非在作业提交阶段校验。OceanBase **Oracle 兼容模式**不支持 `table_options`,配置非空时任务启动即失败。
+- **PostgreSQL**:`fillfactor` 会生成 `WITH (fillfactor=<n>)`,取值须为 `[10, 100]`
的整数;`tablespace` 会生成 `TABLESPACE "..."`,按配置字面量引用(**不受** `fieldIde`
大小写改写)。空白值,以及 `tablespace` 中的非法字符(例如 `"`)会在作业提交阶段被拒绝。仅接受上述 curated key(不支持任意
`WITH` 参数)。OpenGauss、HighGo 通过 Postgres catalog/dialect 继承同一套校验与 DDL。
-SeaTunnel 在提交时仅校验 **key 白名单**,不校验具体取值是否被目标库支持(与 MySQL 首批实现一致)。
+SeaTunnel 在提交时会对所有支持 `table_options` 的方言校验 **key 白名单**。对 PostgreSQL(以及
OpenGauss / HighGo 同源路径),还会校验空白值与 `fillfactor` 数值区间。其他方言(例如
MySQL)在白名单之外不额外校验具体取值是否被目标库支持。
示例(MySQL 自动建表时指定存储引擎与字符集):
@@ -295,6 +297,28 @@ sink {
生成的 DDL 会追加 `ENGINE`、`DEFAULT CHARSET`、`COLLATE` 子句。未在白名单内的 key(如
`bucket_num`)会在作业提交阶段报错。
+示例(PostgreSQL 自动建表时指定 tablespace 与 fillfactor):
+
+```hocon
+sink {
+ Jdbc {
+ url = "jdbc:postgresql://localhost:5432/mydb"
+ driver = "org.postgresql.Driver"
+ username = "postgres"
+ password = "password"
+ database = "mydb"
+ table = "public.orders"
+ generate_sink_sql = true
+ schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST"
+ primary_keys = ["id"]
+ table_options = {
+ "tablespace" = "pg_default"
+ "fillfactor" = "70"
+ }
+ }
+}
+```
+
### custom_sql [String]
当`data_save_mode`选择`CUSTOM_PROCESSING`时,需要填写`CUSTOM_SQL`参数。该参数通常填写一条可以执行的SQL。SQL将在同步任务之前执行
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCatalog.java
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCatalog.java
index fd07f04029..d352f67c8e 100644
---
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCatalog.java
+++
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCatalog.java
@@ -39,6 +39,9 @@ import java.sql.SQLException;
@Slf4j
public class PostgresCatalog extends AbstractJdbcCatalog {
+ public static final String TABLE_OPTION_TABLESPACE = "tablespace";
+ public static final String TABLE_OPTION_FILLFACTOR = "fillfactor";
+
private static final String SELECT_COLUMNS_SQL_TEMPLATE =
"SELECT \n"
+ " a.attname AS column_name, \n"
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCreateTableSqlBuilder.java
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCreateTableSqlBuilder.java
index 3c9528f185..5a2d23cfff 100644
---
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCreateTableSqlBuilder.java
+++
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCreateTableSqlBuilder.java
@@ -44,6 +44,8 @@ public class PostgresCreateTableSqlBuilder {
private PrimaryKey primaryKey;
private String sourceCatalogName;
private String fieldIde;
+ private String tablespace;
+ private String fillfactor;
private List<ConstraintKey> constraintKeys;
public Boolean isHaveConstraintKey = false;
@@ -55,6 +57,8 @@ public class PostgresCreateTableSqlBuilder {
this.primaryKey = catalogTable.getTableSchema().getPrimaryKey();
this.sourceCatalogName = catalogTable.getCatalogName();
this.fieldIde = catalogTable.getOptions().get("fieldIde");
+ this.tablespace =
catalogTable.getOptions().get(PostgresCatalog.TABLE_OPTION_TABLESPACE);
+ this.fillfactor =
catalogTable.getOptions().get(PostgresCatalog.TABLE_OPTION_FILLFACTOR);
this.constraintKeys =
catalogTable.getTableSchema().getConstraintKeys();
this.createIndex = createIndex;
}
@@ -106,7 +110,17 @@ public class PostgresCreateTableSqlBuilder {
}
createTableSql.append(String.join(",\n", columnSqls));
- createTableSql.append("\n);");
+ createTableSql.append("\n)");
+ if (StringUtils.isNotBlank(fillfactor)) {
+ // Value is validated as 10-100 by
PostgresDialect.validateTableOptions.
+ createTableSql.append("\nWITH
(fillfactor=").append(fillfactor.trim()).append(")");
+ }
+ if (StringUtils.isNotBlank(tablespace)) {
+ // Tablespace is a storage object name: quote literally, do NOT
apply fieldIde
+ // case rewriting (that policy is only for table/column
identifiers).
+ createTableSql.append("\nTABLESPACE
\"").append(tablespace.trim()).append("\"");
+ }
+ createTableSql.append(";");
List<String> commentSqls =
columns.stream()
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/psql/PostgresDialect.java
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/psql/PostgresDialect.java
index 365c81971f..ae8087372f 100644
---
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/psql/PostgresDialect.java
+++
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/psql/PostgresDialect.java
@@ -19,6 +19,7 @@ package
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.psql;
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.psql.PostgresCatalog;
+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;
@@ -43,7 +46,9 @@ import java.sql.SQLException;
import java.sql.Statement;
import java.util.ArrayList;
import java.util.Arrays;
+import java.util.Collections;
import java.util.HashSet;
+import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
import java.util.Optional;
@@ -62,6 +67,18 @@ public class PostgresDialect implements JdbcDialect {
private static final long serialVersionUID = -5834746193472465218L;
public static final int DEFAULT_POSTGRES_FETCH_SIZE = 128;
+ /** PostgreSQL FILLFACTOR legal range (inclusive): 10-100. */
+ private static final int FILLFACTOR_MIN = 10;
+
+ private static final int FILLFACTOR_MAX = 100;
+
+ private static final Set<String> SUPPORTED_TABLE_OPTIONS =
+ Collections.unmodifiableSet(
+ new LinkedHashSet<>(
+ Arrays.asList(
+ PostgresCatalog.TABLE_OPTION_TABLESPACE,
+ PostgresCatalog.TABLE_OPTION_FILLFACTOR)));
+
public String fieldIde = FieldIdeEnum.ORIGINAL.getValue();
public PostgresDialect() {}
@@ -462,4 +479,83 @@ public class PostgresDialect implements JdbcDialect {
}
return columnName;
}
+
+ @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 (PostgresCatalog.TABLE_OPTION_FILLFACTOR.equals(key)) {
+ validateFillfactor(trimmed);
+ } else if (PostgresCatalog.TABLE_OPTION_TABLESPACE.equals(key)) {
+ validateTablespace(trimmed);
+ }
+ }
+ }
+
+ private void validateFillfactor(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'",
+ dialectName(),
+ PostgresCatalog.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'",
+ dialectName(),
+ PostgresCatalog.TABLE_OPTION_FILLFACTOR,
+ FILLFACTOR_MIN,
+ FILLFACTOR_MAX,
+ value));
+ }
+ }
+
+ private void validateTablespace(String value) {
+ // Always emitted as TABLESPACE "...", so reject quote / control chars
that break DDL.
+ if (value.indexOf('"') >= 0
+ || value.indexOf('\n') >= 0
+ || value.indexOf('\r') >= 0
+ || value.indexOf(';') >= 0) {
+ throw new JdbcConnectorException(
+ SeaTunnelAPIErrorCode.CONFIG_VALIDATION_FAILED,
+ String.format(
+ "Invalid JDBC table_options for dialect '%s': key
'%s' contains illegal characters: '%s'",
+ dialectName(),
PostgresCatalog.TABLE_OPTION_TABLESPACE, value));
+ }
+ }
}
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCreateTableSqlBuilderTest.java
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCreateTableSqlBuilderTest.java
index b0ddc13e45..e81dea07b9 100644
---
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCreateTableSqlBuilderTest.java
+++
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCreateTableSqlBuilderTest.java
@@ -35,7 +35,9 @@ import org.junit.jupiter.api.Test;
import java.util.Arrays;
import java.util.Collections;
+import java.util.HashMap;
import java.util.List;
+import java.util.Map;
import java.util.regex.Pattern;
import static org.mockito.Mockito.mock;
@@ -88,6 +90,54 @@ class PostgresCreateTableSqlBuilderTest {
});
}
+ @Test
+ void testBuildCreateTableSqlWithTableOptions() {
+ CatalogTable catalogTable = catalogTable(false);
+ Map<String, String> options = new HashMap<>(catalogTable.getOptions());
+ options.put(PostgresCatalog.TABLE_OPTION_TABLESPACE, "pg_default");
+ options.put(PostgresCatalog.TABLE_OPTION_FILLFACTOR, "70");
+ CatalogTable tableWithOptions =
+ CatalogTable.of(
+ catalogTable.getTableId(),
+ catalogTable.getTableSchema(),
+ options,
+ catalogTable.getPartitionKeys(),
+ catalogTable.getComment());
+
+ String createTableSql =
+ new PostgresCreateTableSqlBuilder(tableWithOptions, false)
+ .build(tableWithOptions.getTableId().toTablePath());
+
+ Assertions.assertTrue(createTableSql.contains("WITH (fillfactor=70)"));
+ Assertions.assertTrue(createTableSql.contains("TABLESPACE
\"pg_default\""));
+ }
+
+ @Test
+ void testBuildCreateTableSqlWithTableOptionsIgnoresFieldIde() {
+ CatalogTable catalogTable = catalogTable(false);
+ Map<String, String> options = new HashMap<>(catalogTable.getOptions());
+ options.put("fieldIde", "UPPERCASE");
+ options.put(PostgresCatalog.TABLE_OPTION_TABLESPACE, "pg_default");
+ options.put(PostgresCatalog.TABLE_OPTION_FILLFACTOR, "70");
+ CatalogTable tableWithOptions =
+ CatalogTable.of(
+ catalogTable.getTableId(),
+ catalogTable.getTableSchema(),
+ options,
+ catalogTable.getPartitionKeys(),
+ catalogTable.getComment());
+
+ String createTableSql =
+ new PostgresCreateTableSqlBuilder(tableWithOptions, false)
+ .build(tableWithOptions.getTableId().toTablePath());
+
+ Assertions.assertTrue(createTableSql.contains("WITH (fillfactor=70)"));
+ Assertions.assertTrue(
+ createTableSql.contains("TABLESPACE \"pg_default\""),
+ "tablespace must not be rewritten by fieldIde; got: " +
createTableSql);
+ Assertions.assertFalse(createTableSql.contains("TABLESPACE
\"PG_DEFAULT\""));
+ }
+
private CatalogTable catalogTable(boolean otherDB) {
TableIdentifier tableIdentifier =
TableIdentifier.of(
diff --git
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/psql/PostgresDialectTableOptionsTest.java
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/psql/PostgresDialectTableOptionsTest.java
new file mode 100644
index 0000000000..a26c49fcdb
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/psql/PostgresDialectTableOptionsTest.java
@@ -0,0 +1,126 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.psql;
+
+import
org.apache.seatunnel.connectors.seatunnel.jdbc.exception.JdbcConnectorException;
+
+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 PostgresDialectTableOptionsTest {
+
+ @Test
+ public void testValidateTableOptions() {
+ PostgresDialect dialect = new PostgresDialect();
+ Map<String, String> tableOptions = new HashMap<>();
+ tableOptions.put("tablespace", "pg_default");
+ tableOptions.put("fillfactor", "70");
+
+ Assertions.assertDoesNotThrow(() ->
dialect.validateTableOptions(tableOptions));
+ }
+
+ @Test
+ public void testValidateTableOptionsFillfactorBoundary() {
+ PostgresDialect dialect = new PostgresDialect();
+ Assertions.assertDoesNotThrow(
+ () ->
dialect.validateTableOptions(Collections.singletonMap("fillfactor", "10")));
+ Assertions.assertDoesNotThrow(
+ () ->
dialect.validateTableOptions(Collections.singletonMap("fillfactor", "100")));
+ }
+
+ @Test
+ public void testValidateTableOptionsWithUnknownKey() {
+ PostgresDialect dialect = new PostgresDialect();
+
+ JdbcConnectorException exception =
+ Assertions.assertThrows(
+ JdbcConnectorException.class,
+ () ->
+ dialect.validateTableOptions(
+ Collections.singletonMap("engine",
"InnoDB")));
+ Assertions.assertTrue(exception.getMessage().contains("Unsupported
JDBC table_options"));
+ Assertions.assertTrue(exception.getMessage().contains("Postgres"));
+ }
+
+ @Test
+ public void testValidateTableOptionsRejectBlankValues() {
+ PostgresDialect dialect = new PostgresDialect();
+
+ 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
+ public void testValidateTableOptionsRejectInvalidFillfactor() {
+ PostgresDialect dialect = new PostgresDialect();
+
+ JdbcConnectorException nonNumeric =
+ Assertions.assertThrows(
+ JdbcConnectorException.class,
+ () ->
+ dialect.validateTableOptions(
+ Collections.singletonMap("fillfactor",
"abc")));
+ Assertions.assertTrue(nonNumeric.getMessage().contains("must be an
integer between"));
+
+ JdbcConnectorException tooLow =
+ Assertions.assertThrows(
+ JdbcConnectorException.class,
+ () ->
+ dialect.validateTableOptions(
+ Collections.singletonMap("fillfactor",
"9")));
+ Assertions.assertTrue(tooLow.getMessage().contains("must be an integer
between"));
+
+ JdbcConnectorException tooHigh =
+ Assertions.assertThrows(
+ JdbcConnectorException.class,
+ () ->
+ dialect.validateTableOptions(
+ Collections.singletonMap("fillfactor",
"101")));
+ Assertions.assertTrue(tooHigh.getMessage().contains("must be an
integer between"));
+ }
+
+ @Test
+ public void testValidateTableOptionsRejectIllegalTablespace() {
+ PostgresDialect dialect = new PostgresDialect();
+
+ JdbcConnectorException exception =
+ Assertions.assertThrows(
+ JdbcConnectorException.class,
+ () ->
+ dialect.validateTableOptions(
+ Collections.singletonMap("tablespace",
"pg_\"default\"")));
+ Assertions.assertTrue(exception.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 65d4c31e00..d8d03a9303 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
@@ -48,17 +48,28 @@ class JdbcTableOptionsConditionExtensionTest {
}
@Test
- void testPostgresRejectsNonEmptyTableOptionsViaOptionRule() {
+ void testPostgresTableOptionsPassViaOptionRule() {
Map<String, Object> config = postgresSinkConfig();
Map<String, String> tableOptions = new HashMap<>();
+ tableOptions.put("tablespace", "pg_default");
tableOptions.put("fillfactor", "70");
config.put(SinkConnectorCommonOptions.TABLE_OPTIONS.key(),
tableOptions);
+ Assertions.assertDoesNotThrow(() -> validateSinkOptionRule(config));
+ }
+
+ @Test
+ void testPostgresRejectsUnknownTableOptionsViaOptionRule() {
+ Map<String, Object> config = postgresSinkConfig();
+ 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("not supported for dialect
'Postgres'"));
+ Assertions.assertTrue(exception.getMessage().contains("Unsupported
JDBC table_options"));
+ Assertions.assertTrue(exception.getMessage().contains("Postgres"));
}
@Test