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 30474f8632 [Improve][Connector-V2] Validate OpenMLDB SQL (#12181)
30474f8632 is described below
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());