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-12181-fdb6beb3ff59602f43e85cd6713b4daeb204720a in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit 30474f8632a88b7cc7116f0abd8e8eb83ad80f0a Author: Nikhil kumar <[email protected]> AuthorDate: Tue Sep 8 05:12:15 2026 +0000 [Improve][Connector-V2] Validate OpenMLDB SQL (#12181) --- docs/en/connectors/source/OpenMldb.md | 2 + docs/zh/connectors/source/OpenMldb.md | 2 + .../openmldb/source/OpenMldbSourceFactory.java | 4 +- .../seatunnel/openmldb/OpenMldbFactoryTest.java | 72 ++++++++++++++++++++++ 4 files changed, 79 insertions(+), 1 deletion(-) diff --git a/docs/en/connectors/source/OpenMldb.md b/docs/en/connectors/source/OpenMldb.md index 2881eadcc6..1487c0994c 100644 --- a/docs/en/connectors/source/OpenMldb.md +++ b/docs/en/connectors/source/OpenMldb.md @@ -64,6 +64,8 @@ When it is `true`, configure `zk_host` and `zk_path`. ### sql [string] +The required `sql` must not be empty or whitespace-only in either standalone or cluster mode. + The SQL statement to execute against OpenMLDB. The result set columns become the schema of the emitted SeaTunnel rows. diff --git a/docs/zh/connectors/source/OpenMldb.md b/docs/zh/connectors/source/OpenMldb.md index 03143227a3..f36349aac3 100644 --- a/docs/zh/connectors/source/OpenMldb.md +++ b/docs/zh/connectors/source/OpenMldb.md @@ -62,6 +62,8 @@ OpenMLDB 类型会按照所配置 `sql` 语句的结果集映射为 SeaTunnel ### sql [string] +无论使用单机模式还是集群模式,必填项 `sql` 都不能为空字符串或仅包含空白字符。 + 针对 OpenMLDB 执行的 SQL 语句,结果集的列会成为连接器输出行的字段。 ### database [string] diff --git a/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbSourceFactory.java b/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbSourceFactory.java index def0c84181..0ff80bbac2 100644 --- a/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbSourceFactory.java +++ b/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbSourceFactory.java @@ -31,6 +31,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 OpenMldbSourceFactory implements TableSourceFactory { @Override @@ -42,7 +44,7 @@ public class OpenMldbSourceFactory implements TableSourceFactory { public OptionRule optionRule() { return OptionRule.builder() .required(OpenMldbSourceOptions.CLUSTER_MODE) - .required(OpenMldbSourceOptions.SQL) + .required(OpenMldbSourceOptions.SQL, notBlank(OpenMldbSourceOptions.SQL)) .required(OpenMldbSourceOptions.DATABASE) .optional(OpenMldbSourceOptions.SESSION_TIMEOUT) .optional(OpenMldbSourceOptions.REQUEST_TIMEOUT) diff --git a/seatunnel-connectors-v2/connector-openmldb/src/test/java/org/apache/seatunnel/connectors/seatunnel/openmldb/OpenMldbFactoryTest.java b/seatunnel-connectors-v2/connector-openmldb/src/test/java/org/apache/seatunnel/connectors/seatunnel/openmldb/OpenMldbFactoryTest.java index 89982ce78d..daccae7527 100644 --- a/seatunnel-connectors-v2/connector-openmldb/src/test/java/org/apache/seatunnel/connectors/seatunnel/openmldb/OpenMldbFactoryTest.java +++ b/seatunnel-connectors-v2/connector-openmldb/src/test/java/org/apache/seatunnel/connectors/seatunnel/openmldb/OpenMldbFactoryTest.java @@ -17,13 +17,85 @@ package org.apache.seatunnel.connectors.seatunnel.openmldb; +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.connectors.seatunnel.openmldb.source.OpenMldbSourceFactory; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import java.util.HashMap; +import java.util.Map; class OpenMldbFactoryTest { + private final OptionRule optionRule = new OpenMldbSourceFactory().optionRule(); + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testNonblankSqlAccepted(boolean clusterMode) { + for (String sql : + new String[] { + "select * from test_table", " select * from test_table ", "not parsed here" + }) { + Map<String, Object> config = requiredConfig(clusterMode); + config.put("sql", sql); + Assertions.assertDoesNotThrow(() -> validate(config)); + } + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testMissingSqlRejected(boolean clusterMode) { + Map<String, Object> config = requiredConfig(clusterMode); + config.remove("sql"); + Assertions.assertThrows(OptionValidationException.class, () -> validate(config)); + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testEmptySqlRejected(boolean clusterMode) { + assertInvalidSql(clusterMode, ""); + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testWhitespaceOnlySqlRejected(boolean clusterMode) { + assertInvalidSql(clusterMode, " "); + assertInvalidSql(clusterMode, "\t\r\n"); + } + + private void assertInvalidSql(boolean clusterMode, String sql) { + Map<String, Object> config = requiredConfig(clusterMode); + config.put("sql", sql); + Assertions.assertThrows(OptionValidationException.class, () -> validate(config)); + } + + private void validate(Map<String, Object> config) { + ReadonlyConfig readonlyConfig = ReadonlyConfig.fromMap(config); + ConfigValidator.validateUnknownKeys(readonlyConfig, optionRule, "OpenMldb"); + ConfigValidator.of(readonlyConfig).validate(optionRule); + } + + private Map<String, Object> requiredConfig(boolean clusterMode) { + Map<String, Object> config = new HashMap<>(); + config.put("cluster_mode", clusterMode); + config.put("database", "test_db"); + config.put("sql", "select * from test_table"); + if (clusterMode) { + config.put("zk_host", "localhost:2181"); + config.put("zk_path", "/openmldb"); + } else { + config.put("host", "localhost"); + config.put("port", 6527); + } + return config; + } + @Test void optionRule() { Assertions.assertNotNull((new OpenMldbSourceFactory()).optionRule());
