This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] 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 0e1a8b6f84 Add declarative required-option validation for TDengine
(#12490)
0e1a8b6f84 is described below
commit 0e1a8b6f848e9cdf1dae4881d70d61103fd13cff
Author: Yiğitcan Öztürk <[email protected]>
AuthorDate: Sun Sep 27 11:06:36 2026 +0000
Add declarative required-option validation for TDengine (#12490)
---
docs/en/connectors/sink/TDengine.md | 2 +
docs/en/connectors/source/TDengine.md | 2 +
docs/zh/connectors/sink/TDengine.md | 2 +
docs/zh/connectors/source/TDengine.md | 2 +
.../tdengine/sink/TDengineSinkFactory.java | 13 ++--
.../tdengine/source/TDengineSourceFactory.java | 17 +++--
.../tdengine/sink/TDengineSinkFactoryTest.java | 78 ++++++++++++++++++++
.../tdengine/source/TDengineSourceFactoryTest.java | 82 ++++++++++++++++++++++
8 files changed, 186 insertions(+), 12 deletions(-)
diff --git a/docs/en/connectors/sink/TDengine.md
b/docs/en/connectors/sink/TDengine.md
index cf32eac44a..ca2937b75f 100644
--- a/docs/en/connectors/sink/TDengine.md
+++ b/docs/en/connectors/sink/TDengine.md
@@ -98,6 +98,8 @@ Sink plugin common parameters, please refer to
For multi-table writes, `multi_table_sink_replica` can be used with the common
sink options.
+All required string options listed above must contain a nonblank value (not
empty or whitespace-only). This validation does not check server connectivity
or timestamp format.
+
## Input Row Shape
The connector expects every input row to follow the super-table write shape:
diff --git a/docs/en/connectors/source/TDengine.md
b/docs/en/connectors/source/TDengine.md
index f01d8f394e..fba1913266 100644
--- a/docs/en/connectors/source/TDengine.md
+++ b/docs/en/connectors/source/TDengine.md
@@ -109,6 +109,8 @@ correctly.
Source plugin common parameters, please refer to
[Source Common Options](../common-options/source-common-options.md) for
details.
+All required string options listed above must contain a nonblank value (not
empty or whitespace-only). This validation does not check server connectivity
or timestamp format.
+
## Output Schema
The output table always starts with the reserved `subtable_name` column, which
diff --git a/docs/zh/connectors/sink/TDengine.md
b/docs/zh/connectors/sink/TDengine.md
index 74e2f126a6..775d57a9d7 100644
--- a/docs/zh/connectors/sink/TDengine.md
+++ b/docs/zh/connectors/sink/TDengine.md
@@ -79,6 +79,8 @@ TDengine 服务端时区,用于时间戳转换,默认值为 `UTC`。如果
Sink 插件通用参数,请参考 [Sink Common
Options](../common-options/sink-common-options.md)。
多表写入时,可以配合通用参数中的 `multi_table_sink_replica` 使用。
+上述必填字符串选项均不能为空字符串或仅包含空白字符。此项校验不检查服务器连接或时间戳格式。
+
## 输入数据格式
连接器要求每行输入数据符合超级表写入结构:
diff --git a/docs/zh/connectors/source/TDengine.md
b/docs/zh/connectors/source/TDengine.md
index 3a305d98af..a18a77f929 100644
--- a/docs/zh/connectors/source/TDengine.md
+++ b/docs/zh/connectors/source/TDengine.md
@@ -90,6 +90,8 @@ TDengine 子表名称列表。不配置时读取该超级表下的所有子表
Source 插件通用参数,请参考 [Source Common
Options](../common-options/source-common-options.md)。
+上述必填字符串选项均不能为空字符串或仅包含空白字符。此项校验不检查服务器连接或时间戳格式。
+
## 输出 Schema
输出表的第一列固定为预留字段 `subtable_name`,用来标识该行来自哪个 TDengine 子表。后续字段为 `read_columns`
中声明的列(未设置时为所有列),顺序与 `read_columns` 保持一致;TAGS 列请按声明的顺序排在普通列之后。
diff --git
a/seatunnel-connectors-v2/connector-tdengine/src/main/java/org/apache/seatunnel/connectors/seatunnel/tdengine/sink/TDengineSinkFactory.java
b/seatunnel-connectors-v2/connector-tdengine/src/main/java/org/apache/seatunnel/connectors/seatunnel/tdengine/sink/TDengineSinkFactory.java
index 976f5a8728..9c93536fca 100644
---
a/seatunnel-connectors-v2/connector-tdengine/src/main/java/org/apache/seatunnel/connectors/seatunnel/tdengine/sink/TDengineSinkFactory.java
+++
b/seatunnel-connectors-v2/connector-tdengine/src/main/java/org/apache/seatunnel/connectors/seatunnel/tdengine/sink/TDengineSinkFactory.java
@@ -28,6 +28,8 @@ import
org.apache.seatunnel.connectors.seatunnel.tdengine.config.TDengineSinkOpt
import com.google.auto.service.AutoService;
+import static org.apache.seatunnel.api.configuration.util.Conditions.notBlank;
+
@AutoService(Factory.class)
public class TDengineSinkFactory implements TableSinkFactory {
@Override
@@ -38,12 +40,11 @@ public class TDengineSinkFactory implements
TableSinkFactory {
@Override
public OptionRule optionRule() {
return OptionRule.builder()
- .required(
- TDengineSinkOptions.URL,
- TDengineSinkOptions.USERNAME,
- TDengineSinkOptions.PASSWORD,
- TDengineSinkOptions.DATABASE,
- TDengineSinkOptions.STABLE)
+ .required(TDengineSinkOptions.URL,
notBlank(TDengineSinkOptions.URL))
+ .required(TDengineSinkOptions.USERNAME,
notBlank(TDengineSinkOptions.USERNAME))
+ .required(TDengineSinkOptions.PASSWORD,
notBlank(TDengineSinkOptions.PASSWORD))
+ .required(TDengineSinkOptions.DATABASE,
notBlank(TDengineSinkOptions.DATABASE))
+ .required(TDengineSinkOptions.STABLE,
notBlank(TDengineSinkOptions.STABLE))
.optional(
TDengineSinkOptions.TIMEZONE,
SinkConnectorCommonOptions.MULTI_TABLE_SINK_REPLICA)
diff --git
a/seatunnel-connectors-v2/connector-tdengine/src/main/java/org/apache/seatunnel/connectors/seatunnel/tdengine/source/TDengineSourceFactory.java
b/seatunnel-connectors-v2/connector-tdengine/src/main/java/org/apache/seatunnel/connectors/seatunnel/tdengine/source/TDengineSourceFactory.java
index c875593ef0..c0e643e915 100644
---
a/seatunnel-connectors-v2/connector-tdengine/src/main/java/org/apache/seatunnel/connectors/seatunnel/tdengine/source/TDengineSourceFactory.java
+++
b/seatunnel-connectors-v2/connector-tdengine/src/main/java/org/apache/seatunnel/connectors/seatunnel/tdengine/source/TDengineSourceFactory.java
@@ -30,6 +30,8 @@ import com.google.auto.service.AutoService;
import java.io.Serializable;
+import static org.apache.seatunnel.api.configuration.util.Conditions.notBlank;
+
@AutoService(Factory.class)
public class TDengineSourceFactory implements TableSourceFactory {
@@ -41,14 +43,17 @@ public class TDengineSourceFactory implements
TableSourceFactory {
@Override
public OptionRule optionRule() {
return OptionRule.builder()
+ .required(TDengineSourceOptions.URL,
notBlank(TDengineSourceOptions.URL))
+ .required(TDengineSourceOptions.USERNAME,
notBlank(TDengineSourceOptions.USERNAME))
+ .required(TDengineSourceOptions.PASSWORD,
notBlank(TDengineSourceOptions.PASSWORD))
+ .required(TDengineSourceOptions.DATABASE,
notBlank(TDengineSourceOptions.DATABASE))
+ .required(TDengineSourceOptions.STABLE,
notBlank(TDengineSourceOptions.STABLE))
.required(
- TDengineSourceOptions.URL,
- TDengineSourceOptions.USERNAME,
- TDengineSourceOptions.PASSWORD,
- TDengineSourceOptions.DATABASE,
- TDengineSourceOptions.STABLE,
TDengineSourceOptions.LOWER_BOUND,
- TDengineSourceOptions.UPPER_BOUND)
+ notBlank(TDengineSourceOptions.LOWER_BOUND))
+ .required(
+ TDengineSourceOptions.UPPER_BOUND,
+ notBlank(TDengineSourceOptions.UPPER_BOUND))
.build();
}
diff --git
a/seatunnel-connectors-v2/connector-tdengine/src/test/java/org/apache/seatunnel/connectors/seatunnel/tdengine/sink/TDengineSinkFactoryTest.java
b/seatunnel-connectors-v2/connector-tdengine/src/test/java/org/apache/seatunnel/connectors/seatunnel/tdengine/sink/TDengineSinkFactoryTest.java
new file mode 100644
index 0000000000..daa0049c79
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-tdengine/src/test/java/org/apache/seatunnel/connectors/seatunnel/tdengine/sink/TDengineSinkFactoryTest.java
@@ -0,0 +1,78 @@
+/*
+ * 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.tdengine.sink;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConfigValidator;
+import org.apache.seatunnel.api.configuration.util.OptionValidationException;
+import
org.apache.seatunnel.connectors.seatunnel.tdengine.config.TDengineSinkOptions;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.HashMap;
+import java.util.Map;
+
+class TDengineSinkFactoryTest {
+
+ @Test
+ void validRequiredOptions() {
+ Assertions.assertDoesNotThrow(() -> validate(validConfig()));
+ }
+
+ @Test
+ void requiredOptionsRejectMissingEmptyAndWhitespace() {
+ for (String key : requiredKeys()) {
+ Map<String, Object> missing = validConfig();
+ missing.remove(key);
+ Assertions.assertThrows(OptionValidationException.class, () ->
validate(missing), key);
+
+ for (String invalid : new String[] {"", " "}) {
+ Map<String, Object> config = validConfig();
+ config.put(key, invalid);
+ Assertions.assertThrows(
+ OptionValidationException.class, () ->
validate(config), key);
+ }
+ }
+ }
+
+ private void validate(Map<String, Object> config) {
+ ConfigValidator.of(ReadonlyConfig.fromMap(config))
+ .validate(new TDengineSinkFactory().optionRule());
+ }
+
+ private String[] requiredKeys() {
+ return new String[] {
+ TDengineSinkOptions.URL.key(),
+ TDengineSinkOptions.USERNAME.key(),
+ TDengineSinkOptions.PASSWORD.key(),
+ TDengineSinkOptions.DATABASE.key(),
+ TDengineSinkOptions.STABLE.key()
+ };
+ }
+
+ private Map<String, Object> validConfig() {
+ Map<String, Object> config = new HashMap<>();
+ config.put(TDengineSinkOptions.URL.key(),
"jdbc:TAOS-RS://localhost:6041");
+ config.put(TDengineSinkOptions.USERNAME.key(), "username");
+ config.put(TDengineSinkOptions.PASSWORD.key(), "password");
+ config.put(TDengineSinkOptions.DATABASE.key(), "database");
+ config.put(TDengineSinkOptions.STABLE.key(), "stable");
+ return config;
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-tdengine/src/test/java/org/apache/seatunnel/connectors/seatunnel/tdengine/source/TDengineSourceFactoryTest.java
b/seatunnel-connectors-v2/connector-tdengine/src/test/java/org/apache/seatunnel/connectors/seatunnel/tdengine/source/TDengineSourceFactoryTest.java
new file mode 100644
index 0000000000..df4a5eb32a
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-tdengine/src/test/java/org/apache/seatunnel/connectors/seatunnel/tdengine/source/TDengineSourceFactoryTest.java
@@ -0,0 +1,82 @@
+/*
+ * 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.tdengine.source;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConfigValidator;
+import org.apache.seatunnel.api.configuration.util.OptionValidationException;
+import
org.apache.seatunnel.connectors.seatunnel.tdengine.config.TDengineSourceOptions;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.HashMap;
+import java.util.Map;
+
+class TDengineSourceFactoryTest {
+
+ @Test
+ void validRequiredOptions() {
+ Assertions.assertDoesNotThrow(() -> validate(validConfig()));
+ }
+
+ @Test
+ void requiredOptionsRejectMissingEmptyAndWhitespace() {
+ for (String key : requiredKeys()) {
+ Map<String, Object> missing = validConfig();
+ missing.remove(key);
+ Assertions.assertThrows(OptionValidationException.class, () ->
validate(missing), key);
+
+ for (String invalid : new String[] {"", " "}) {
+ Map<String, Object> config = validConfig();
+ config.put(key, invalid);
+ Assertions.assertThrows(
+ OptionValidationException.class, () ->
validate(config), key);
+ }
+ }
+ }
+
+ private void validate(Map<String, Object> config) {
+ ConfigValidator.of(ReadonlyConfig.fromMap(config))
+ .validate(new TDengineSourceFactory().optionRule());
+ }
+
+ private String[] requiredKeys() {
+ return new String[] {
+ TDengineSourceOptions.URL.key(),
+ TDengineSourceOptions.USERNAME.key(),
+ TDengineSourceOptions.PASSWORD.key(),
+ TDengineSourceOptions.DATABASE.key(),
+ TDengineSourceOptions.STABLE.key(),
+ TDengineSourceOptions.LOWER_BOUND.key(),
+ TDengineSourceOptions.UPPER_BOUND.key()
+ };
+ }
+
+ private Map<String, Object> validConfig() {
+ Map<String, Object> config = new HashMap<>();
+ config.put(TDengineSourceOptions.URL.key(),
"jdbc:TAOS-RS://localhost:6041");
+ config.put(TDengineSourceOptions.USERNAME.key(), "username");
+ config.put(TDengineSourceOptions.PASSWORD.key(), "password");
+ config.put(TDengineSourceOptions.DATABASE.key(), "database");
+ config.put(TDengineSourceOptions.STABLE.key(), "stable");
+ config.put(TDengineSourceOptions.LOWER_BOUND.key(), "2020-01-01
00:00:00");
+ config.put(TDengineSourceOptions.UPPER_BOUND.key(), "2020-01-02
00:00:00");
+ return config;
+ }
+}