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 6bb02da900 [Improve][Connector-V2][Iceberg] Align OptionRule with
docs and strengthen factory test (#11675)
6bb02da900 is described below
commit 6bb02da900f40a540a5c8caf211d8c579e124234
Author: Claire <[email protected]>
AuthorDate: Mon Aug 17 16:34:48 2026 +0800
[Improve][Connector-V2][Iceberg] Align OptionRule with docs and strengthen
factory test (#11675)
---
docs/en/connectors/sink/Iceberg.md | 6 +-
docs/en/connectors/source/Iceberg.md | 4 +-
docs/zh/connectors/sink/Iceberg.md | 6 +-
docs/zh/connectors/source/Iceberg.md | 4 +-
.../seatunnel/iceberg/sink/IcebergSinkFactory.java | 5 +-
.../iceberg/source/IcebergSourceFactory.java | 8 +-
.../seatunnel/iceberg/IcebergFactoryTest.java | 182 ++++++++++++++++++++-
7 files changed, 196 insertions(+), 19 deletions(-)
diff --git a/docs/en/connectors/sink/Iceberg.md
b/docs/en/connectors/sink/Iceberg.md
index 10ce5bc050..26d3395b58 100644
--- a/docs/en/connectors/sink/Iceberg.md
+++ b/docs/en/connectors/sink/Iceberg.md
@@ -64,9 +64,9 @@ libfb303-xxx.jar
| Name | Type | Required | Default
| Description
|
|----------------------------------------|---------|----------|------------------------------|---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|
-| catalog_name | string | yes | default
| User-specified catalog name. default is `default`
|
-| namespace | string | yes | default
| The iceberg database name in the backend catalog. default is
`default`
|
-| table | string | yes | -
| The iceberg table name in the backend catalog.
|
+| catalog_name | string | no | default
| User-specified catalog name. default is `default`
|
+| namespace | string | no | default
| The iceberg database name in the backend catalog. default is
`default`
|
+| table | string | no | -
| The iceberg table name in the backend catalog. If not set, the
table name of the upstream table is used.
|
| iceberg.catalog.config | map | yes | -
| Specify the properties for initializing the Iceberg catalog,
which can be referenced in this file:
[CatalogProperties.java](https://github.com/apache/iceberg/blob/main/core/src/main/java/org/apache/iceberg/CatalogProperties.java)
|
| hadoop.config | map | no | -
| Properties passed through to the Hadoop configuration
|
| iceberg.hadoop-conf-path | string | no | -
| The specified loading paths for the 'core-site.xml',
'hdfs-site.xml', 'hive-site.xml' files.
|
diff --git a/docs/en/connectors/source/Iceberg.md
b/docs/en/connectors/source/Iceberg.md
index 8765e4cf2c..338bbe1936 100644
--- a/docs/en/connectors/source/Iceberg.md
+++ b/docs/en/connectors/source/Iceberg.md
@@ -75,8 +75,8 @@ libfb303-xxx.jar
| Name | Type | Required | Default |
Description
[...]
|--------------------------|---------|----------|----------------------|------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
[...]
-| catalog_name | string | yes | - |
User-specified catalog name.
[...]
-| namespace | string | yes | - | The
iceberg database name in the backend catalog.
[...]
+| catalog_name | string | no | default |
User-specified catalog name.
[...]
+| namespace | string | no | default | The
iceberg database name in the backend catalog.
[...]
| table | string | no | - | The
Iceberg table name in the backend catalog. Configure exactly one of `table` and
`table_list`.
[...]
| table_list | list | no | - | The
Iceberg table list in the backend catalog. Configure exactly one of `table` and
`table_list`. Each item can set `table`, `query`, and snapshot/stream scan
options for that table.
[...]
| iceberg.catalog.config | map | yes | - |
Specify the properties for initializing the Iceberg catalog, which can be
referenced in this file:
[CatalogProperties.java](https://github.com/apache/iceberg/blob/main/core/src/main/java/org/apache/iceberg/CatalogProperties.java)
[...]
diff --git a/docs/zh/connectors/sink/Iceberg.md
b/docs/zh/connectors/sink/Iceberg.md
index e98a0b4ca4..ad3e7d35fb 100644
--- a/docs/zh/connectors/sink/Iceberg.md
+++ b/docs/zh/connectors/sink/Iceberg.md
@@ -64,9 +64,9 @@ libfb303-xxx.jar
| 名称 | 类型 | 是否必须 | 默认
| 描述
|
|----------------------------------------|---------|------|------------------------------|-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|
-| catalog_name | string | yes | default
| 用户指定的目录名称,默认为`default`
|
-| namespace | string | yes | default
| backend catalog(元数据存储的后端目录)中 Iceberg 数据库的名称,默认为 `default`
|
-| table | string | yes | -
| backend catalog(元数据存储的后端目录)中 Iceberg 表的名称
|
+| catalog_name | string | no | default
| 用户指定的目录名称,默认为`default`
|
+| namespace | string | no | default
| backend catalog(元数据存储的后端目录)中 Iceberg 数据库的名称,默认为 `default`
|
+| table | string | no | -
| backend catalog(元数据存储的后端目录)中 Iceberg 表的名称。不配置时使用上游表的表名
|
| iceberg.catalog.config | map | yes | -
| 用于指定初始化 Iceberg Catalog
的属性,这些属性可以参考此文件:[CatalogProperties.java](https://github.com/apache/iceberg/blob/main/core/src/main/java/org/apache/iceberg/CatalogProperties.java)
|
| hadoop.config | map | no | -
| 传递给 Hadoop 配置的属性
|
| iceberg.hadoop-conf-path | string | no | -
| 指定`core-site.xml`、`hdfs-site.xml`、`hive-site.xml` 文件的加载路径
|
diff --git a/docs/zh/connectors/source/Iceberg.md
b/docs/zh/connectors/source/Iceberg.md
index 41a9617419..11797afc85 100644
--- a/docs/zh/connectors/source/Iceberg.md
+++ b/docs/zh/connectors/source/Iceberg.md
@@ -75,8 +75,8 @@ libfb303-xxx.jar
| 参数名 | 类型 | 必须 | 默认值 | 描述
[...]
|--------------------------|---------|------|----------------------|----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
[...]
-| catalog_name | string | 是 | - | 用户指定的目录名称。
[...]
-| namespace | string | 是 | - | 后端目录中的
iceberg 数据库名称。
[...]
+| catalog_name | string | 否 | default | 用户指定的目录名称。
[...]
+| namespace | string | 否 | default | 后端目录中的
iceberg 数据库名称。
[...]
| table | string | 否 | - | 后端目录中的
Iceberg 表名称。`table` 和 `table_list` 必须二选一配置,不能同时配置。
[...]
| table_list | list | 否 | - | 后端目录中的
Iceberg 表列表。`table` 和 `table_list` 必须二选一配置,不能同时配置。每一项都可以配置 `table`、`query`
以及该表自己的快照/流式扫描参数。
[...]
| iceberg.catalog.config | map | 是 | - | 指定初始化
Iceberg
目录的属性,可以在此文件中引用:[CatalogProperties.java](https://github.com/apache/iceberg/blob/main/core/src/main/java/org/apache/iceberg/CatalogProperties.java)
[...]
diff --git
a/seatunnel-connectors-v2/connector-iceberg/src/main/java/org/apache/seatunnel/connectors/seatunnel/iceberg/sink/IcebergSinkFactory.java
b/seatunnel-connectors-v2/connector-iceberg/src/main/java/org/apache/seatunnel/connectors/seatunnel/iceberg/sink/IcebergSinkFactory.java
index 4fec73f442..f71aa440db 100644
---
a/seatunnel-connectors-v2/connector-iceberg/src/main/java/org/apache/seatunnel/connectors/seatunnel/iceberg/sink/IcebergSinkFactory.java
+++
b/seatunnel-connectors-v2/connector-iceberg/src/main/java/org/apache/seatunnel/connectors/seatunnel/iceberg/sink/IcebergSinkFactory.java
@@ -46,12 +46,11 @@ public class IcebergSinkFactory implements TableSinkFactory
{
@Override
public OptionRule optionRule() {
return OptionRule.builder()
- .required(
+ .required(IcebergSinkOptions.CATALOG_PROPS)
+ .optional(
IcebergCommonOptions.KEY_CATALOG_NAME,
IcebergSinkOptions.KEY_NAMESPACE,
IcebergSinkOptions.KEY_TABLE,
- IcebergSinkOptions.CATALOG_PROPS)
- .optional(
IcebergSinkOptions.HADOOP_PROPS,
IcebergSinkOptions.HADOOP_CONF_PATH_PROP,
IcebergSinkOptions.KEY_CASE_SENSITIVE,
diff --git
a/seatunnel-connectors-v2/connector-iceberg/src/main/java/org/apache/seatunnel/connectors/seatunnel/iceberg/source/IcebergSourceFactory.java
b/seatunnel-connectors-v2/connector-iceberg/src/main/java/org/apache/seatunnel/connectors/seatunnel/iceberg/source/IcebergSourceFactory.java
index b3bf791983..49538f6905 100644
---
a/seatunnel-connectors-v2/connector-iceberg/src/main/java/org/apache/seatunnel/connectors/seatunnel/iceberg/source/IcebergSourceFactory.java
+++
b/seatunnel-connectors-v2/connector-iceberg/src/main/java/org/apache/seatunnel/connectors/seatunnel/iceberg/source/IcebergSourceFactory.java
@@ -56,13 +56,13 @@ public class IcebergSourceFactory implements
TableSourceFactory {
@Override
public OptionRule optionRule() {
return OptionRule.builder()
- .required(
- IcebergCommonOptions.KEY_CATALOG_NAME,
- IcebergCommonOptions.KEY_NAMESPACE,
- IcebergCommonOptions.CATALOG_PROPS)
+ .required(IcebergCommonOptions.CATALOG_PROPS)
.exclusive(IcebergCommonOptions.KEY_TABLE,
IcebergSourceOptions.KEY_TABLE_LIST)
.optional(
+ IcebergCommonOptions.KEY_CATALOG_NAME,
+ IcebergCommonOptions.KEY_NAMESPACE,
ConnectorCommonOptions.SCHEMA,
+ IcebergSourceOptions.QUERY,
IcebergSourceOptions.KEY_CASE_SENSITIVE,
IcebergSourceOptions.KEY_START_SNAPSHOT_TIMESTAMP,
IcebergSourceOptions.KEY_START_SNAPSHOT_ID,
diff --git
a/seatunnel-connectors-v2/connector-iceberg/src/test/java/org/apache/seatunnel/connectors/seatunnel/iceberg/IcebergFactoryTest.java
b/seatunnel-connectors-v2/connector-iceberg/src/test/java/org/apache/seatunnel/connectors/seatunnel/iceberg/IcebergFactoryTest.java
index a4c0753ab8..939dbb9c44 100644
---
a/seatunnel-connectors-v2/connector-iceberg/src/test/java/org/apache/seatunnel/connectors/seatunnel/iceberg/IcebergFactoryTest.java
+++
b/seatunnel-connectors-v2/connector-iceberg/src/test/java/org/apache/seatunnel/connectors/seatunnel/iceberg/IcebergFactoryTest.java
@@ -17,15 +17,193 @@
package org.apache.seatunnel.connectors.seatunnel.iceberg;
+import org.apache.seatunnel.api.configuration.Option;
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConfigValidator;
+import org.apache.seatunnel.api.configuration.util.OptionRule;
+import org.apache.seatunnel.api.configuration.util.OptionValidationException;
+import org.apache.seatunnel.api.configuration.util.RequiredOption;
+import
org.apache.seatunnel.connectors.seatunnel.iceberg.config.IcebergCommonOptions;
+import
org.apache.seatunnel.connectors.seatunnel.iceberg.config.IcebergSinkOptions;
+import
org.apache.seatunnel.connectors.seatunnel.iceberg.config.IcebergSourceOptions;
+import
org.apache.seatunnel.connectors.seatunnel.iceberg.sink.IcebergSinkFactory;
import
org.apache.seatunnel.connectors.seatunnel.iceberg.source.IcebergSourceFactory;
import org.junit.jupiter.api.Assertions;
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.stream.Collectors;
+
class IcebergFactoryTest {
+ private static final OptionRule SOURCE_RULE = new
IcebergSourceFactory().optionRule();
+ private static final OptionRule SINK_RULE = new
IcebergSinkFactory().optionRule();
+
+ // rule structure
+
+ @Test
+ void sourceOptionalContainsQuery() {
+ Assertions.assertTrue(
+
SOURCE_RULE.getOptionalOptions().contains(IcebergSourceOptions.QUERY));
+ }
+
+ @Test
+ void sourceCatalogNameAndNamespaceAreOptional() {
+ Assertions.assertTrue(
+
SOURCE_RULE.getOptionalOptions().contains(IcebergCommonOptions.KEY_CATALOG_NAME));
+ Assertions.assertTrue(
+
SOURCE_RULE.getOptionalOptions().contains(IcebergCommonOptions.KEY_NAMESPACE));
+ List<Option<?>> required = absolutelyRequiredOptions(SOURCE_RULE);
+
Assertions.assertFalse(required.contains(IcebergCommonOptions.KEY_CATALOG_NAME));
+
Assertions.assertFalse(required.contains(IcebergCommonOptions.KEY_NAMESPACE));
+ }
+
+ @Test
+ void sourceCatalogPropsIsRequired() {
+ Assertions.assertTrue(
+ absolutelyRequiredOptions(SOURCE_RULE)
+ .contains(IcebergCommonOptions.CATALOG_PROPS));
+ }
+
+ @Test
+ void sourceTableAndTableListAreExclusive() {
+ boolean hasExclusive =
+ SOURCE_RULE.getRequiredOptions().stream()
+ .filter(o -> o instanceof
RequiredOption.ExclusiveRequiredOptions)
+ .anyMatch(
+ o ->
+
o.getOptions().contains(IcebergCommonOptions.KEY_TABLE)
+ && o.getOptions()
+ .contains(
+
IcebergSourceOptions
+
.KEY_TABLE_LIST));
+ Assertions.assertTrue(hasExclusive);
+ }
+
@Test
- void optionRule() {
- Assertions.assertNotNull((new IcebergSourceFactory()).optionRule());
+ void sinkTableIsOptional() {
+ Assertions.assertTrue(
+
SINK_RULE.getOptionalOptions().contains(IcebergSinkOptions.KEY_TABLE));
+ Assertions.assertFalse(
+
absolutelyRequiredOptions(SINK_RULE).contains(IcebergSinkOptions.KEY_TABLE));
+ }
+
+ @Test
+ void sinkCatalogPropsIsRequired() {
+ Assertions.assertTrue(
+
absolutelyRequiredOptions(SINK_RULE).contains(IcebergSinkOptions.CATALOG_PROPS));
+ }
+
+ // accepted configs
+
+ @Test
+ void sourceMinimalSingleTableValid() {
+ Assertions.assertDoesNotThrow(() ->
validateSource(sourceConfigWithTable()));
+ }
+
+ @Test
+ void sourceValidWithoutCatalogNameAndNamespace() {
+ Map<String, Object> config = sourceConfigWithTable();
+ config.remove("catalog_name");
+ config.remove("namespace");
+ Assertions.assertDoesNotThrow(() -> validateSource(config));
+ }
+
+ @Test
+ void sourceTableListValid() {
+ Map<String, Object> config = baseSourceConfig();
+ Map<String, Object> t1 = new HashMap<>();
+ t1.put("table", "t1");
+ Map<String, Object> t2 = new HashMap<>();
+ t2.put("table", "t2");
+ config.put("table_list", Arrays.asList(t1, t2));
+ Assertions.assertDoesNotThrow(() -> validateSource(config));
+ }
+
+ @Test
+ void sinkValidWithoutTable() {
+ Map<String, Object> config = new HashMap<>();
+ config.put("iceberg.catalog.config", catalogProps());
+ Assertions.assertDoesNotThrow(() -> validateSink(config));
+ }
+
+ // rejected configs
+
+ @Test
+ void sourceMissingCatalogPropsRejected() {
+ Map<String, Object> config = sourceConfigWithTable();
+ config.remove("iceberg.catalog.config");
+ Assertions.assertThrows(OptionValidationException.class, () ->
validateSource(config));
+ }
+
+ @Test
+ void sourceMissingBothTableAndTableListRejected() {
+ Assertions.assertThrows(
+ OptionValidationException.class, () ->
validateSource(baseSourceConfig()));
+ }
+
+ @Test
+ void sourceBothTableAndTableListRejected() {
+ Map<String, Object> config = sourceConfigWithTable();
+ Map<String, Object> t1 = new HashMap<>();
+ t1.put("table", "t1");
+ config.put("table_list", Collections.singletonList(t1));
+ Assertions.assertThrows(OptionValidationException.class, () ->
validateSource(config));
+ }
+
+ @Test
+ void sinkMissingCatalogPropsRejected() {
+ Map<String, Object> config = new HashMap<>();
+ config.put("table", "t1");
+ Assertions.assertThrows(OptionValidationException.class, () ->
validateSink(config));
+ }
+
+ // helpers
+
+ /**
+ * Returns only the options that ConfigValidator treats as unconditionally
required (i.e. {@link
+ * RequiredOption.AbsolutelyRequiredOptions}). Exclusive-group (e.g.
table/table_list) and
+ * conditional requirements are intentionally excluded: they are enforced
by separate validation
+ * paths and are asserted separately in this test.
+ */
+ private static List<Option<?>> absolutelyRequiredOptions(OptionRule rule) {
+ return rule.getRequiredOptions().stream()
+ .filter(o -> o instanceof
RequiredOption.AbsolutelyRequiredOptions)
+ .flatMap(o -> o.getOptions().stream())
+ .collect(Collectors.toList());
+ }
+
+ private static void validateSource(Map<String, Object> config) {
+
ConfigValidator.of(ReadonlyConfig.fromMap(config)).validate(SOURCE_RULE);
+ }
+
+ private static void validateSink(Map<String, Object> config) {
+ ConfigValidator.of(ReadonlyConfig.fromMap(config)).validate(SINK_RULE);
+ }
+
+ private static Map<String, Object> catalogProps() {
+ Map<String, Object> props = new HashMap<>();
+ props.put("type", "hadoop");
+ props.put("warehouse", "file:///tmp/seatunnel/iceberg/");
+ return props;
+ }
+
+ private static Map<String, Object> baseSourceConfig() {
+ Map<String, Object> config = new HashMap<>();
+ config.put("catalog_name", "seatunnel");
+ config.put("namespace", "database1");
+ config.put("iceberg.catalog.config", catalogProps());
+ return config;
+ }
+
+ private static Map<String, Object> sourceConfigWithTable() {
+ Map<String, Object> config = baseSourceConfig();
+ config.put("table", "source_table");
+ return config;
}
}