This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-12490-c3f06f9a82c119e427df0fb1f440bae583148f62 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit 7b94b38bee5d7a93995013f666d32098f7938863 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; + } +}
