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 71860b4469 [Feature][Connector-V2] Support OpenMLDB multi-table source
(#12303)
71860b4469 is described below
commit 71860b4469bbd83722c81aa85ddd7e57f5f440de
Author: Goutam Adwant <[email protected]>
AuthorDate: Sat Sep 19 14:49:58 2026 +0000
[Feature][Connector-V2] Support OpenMLDB multi-table source (#12303)
Signed-off-by: Goutam Adwant <[email protected]>
---
docs/en/connectors/source/OpenMldb.md | 88 ++++-
docs/zh/connectors/source/OpenMldb.md | 81 ++++-
.../openmldb/config/OpenMldbSqlExecutor.java | 71 ++--
.../openmldb/source/OpenMldbReadClient.java | 209 +++++++++++
.../seatunnel/openmldb/source/OpenMldbSource.java | 45 ++-
.../openmldb/source/OpenMldbSourceFactory.java | 125 ++++++-
.../openmldb/source/OpenMldbSourceReader.java | 121 ++++---
.../seatunnel/openmldb/OpenMldbFactoryTest.java | 133 +++++++
.../openmldb/source/OpenMldbSourceReaderTest.java | 238 +++++++++++++
.../openmldb/source/OpenMldbSourceTest.java | 143 ++++++++
.../openmldb/source/TestOpenMldbSourceIT.java | 386 +++++++++++++++++++++
11 files changed, 1526 insertions(+), 114 deletions(-)
diff --git a/docs/en/connectors/source/OpenMldb.md
b/docs/en/connectors/source/OpenMldb.md
index 1487c0994c..4946bcf898 100644
--- a/docs/en/connectors/source/OpenMldb.md
+++ b/docs/en/connectors/source/OpenMldb.md
@@ -16,6 +16,10 @@ Used to read data from OpenMLDB. The connector executes the
configured SQL state
OpenMLDB and turns the result rows into SeaTunnel records. Both standalone and
cluster deployment
modes are supported.
+Queries read the online tables directly. Cluster reads do not submit offline
Spark jobs, return
+job metadata, or change OpenMLDB session/global execution-mode settings.
Reading offline feature
+data is not supported by this connector.
+
## Key features
- [x] [batch](../../introduction/concepts/connector-v2-features.md)
@@ -24,12 +28,14 @@ modes are supported.
- [x] [column projection](../../introduction/concepts/connector-v2-features.md)
- [ ] [parallelism](../../introduction/concepts/connector-v2-features.md)
- [ ] [support user-defined
split](../../introduction/concepts/connector-v2-features.md)
+- [x] [support multiple table
read](../../introduction/concepts/connector-v2-features.md)
## Data Type Mapping
-OpenMLDB types are mapped to SeaTunnel types according to the result schema of
the configured `sql`
-statement. Columns whose types are not natively understood by SeaTunnel will
cause the read to fail with
-an `UNSUPPORTED_DATA_TYPE` error.
+In multi-table mode, `schema.fields` declares the names and types returned by
each query.
+The connector validates the query result against this schema before emitting
rows.
+SQL `NULL` values remain `null`, including nullable numeric and boolean
columns; they are not
+converted to zero or `false`.
| OpenMLDB Data Type | SeaTunnel Data Type |
|--------------------|---------------------|
@@ -47,7 +53,8 @@ an `UNSUPPORTED_DATA_TYPE` error.
| name | type | required | default value | description
|
|-----------------|---------|----------|---------------|----------------------------------------------------------------------------------------|
| cluster_mode | boolean | yes | - | Whether to connect to
OpenMLDB in cluster mode. Set to `false` for standalone mode. |
-| sql | string | yes | - | SQL statement to
execute against OpenMLDB. Column names and types follow the result. |
+| sql | string | conditional | - | Single-query SQL.
Configure either this option or `tables_configs`, not both. |
+| tables_configs | list | conditional | - | Queries and explicit
result schemas for a multi-table read. |
| database | string | yes | - | The OpenMLDB database
name to connect to. |
| host | string | no | - | Required when
`cluster_mode` is `false`. Host of the standalone OpenMLDB server. |
| port | int | no | - | Required when
`cluster_mode` is `false`. Port of the standalone OpenMLDB server. |
@@ -64,10 +71,32 @@ 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.
+When `tables_configs` is absent, `sql` is required and must not be empty or
whitespace-only.
+
+This legacy mode discovers the schema using the SDK's input-schema API. The
result columns must
+match that input schema in count, order and type. Use `tables_configs` with an
explicit result
+schema for projected or aliased query results.
+
+### tables_configs [list]
+
+A non-empty list of queries on the same OpenMLDB instance. Each entry contains:
+
+- `sql`: a non-blank SQL query.
+- `database`: an optional override of the required root-level database.
+- `schema.table`: a unique output table identity for downstream routing.
+- `schema.fields`: all query output column names and their supported SeaTunnel
types.
+
+Field names must exactly match the result column names, including case. Use
SQL aliases where
+needed. Fields are matched by name, so the order of `schema.fields` does not
need to match the
+SQL projection order. Missing, duplicate, extra or incorrectly typed result
columns fail the read.
+
+Keep connection and timeout options at source level. Do not combine
`tables_configs` with
+root-level `sql` or `schema`. Each entry may use a different result schema.
-The SQL statement to execute against OpenMLDB. The result set columns become
the schema of the
-emitted SeaTunnel rows.
+One reader executes the queries sequentially. In batch mode, completion is
signalled only after
+all queries succeed, including empty results. In streaming mode, each poll
executes the queries
+again; this is not CDC or incremental polling and can produce duplicate
records. Multi-table
+reading does not add parallelism, a cross-table consistent snapshot, or
exactly-once guarantees.
### database [string]
@@ -159,6 +188,51 @@ sink {
}
```
+### Multi-table read
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "BATCH"
+}
+
+source {
+ OpenMldb {
+ cluster_mode = false
+ host = "openmldb"
+ port = 6527
+ database = "shop"
+ tables_configs = [
+ {
+ sql = "select id, amount from orders"
+ schema {
+ table = "shop.orders"
+ fields {
+ id = STRING
+ amount = INT
+ }
+ }
+ },
+ {
+ database = "crm"
+ sql = "select id, name from customers"
+ schema {
+ table = "crm.customers"
+ fields {
+ id = STRING
+ name = STRING
+ }
+ }
+ }
+ ]
+ }
+}
+
+sink {
+ Console {}
+}
+```
+
## Changelog
<ChangeLog />
diff --git a/docs/zh/connectors/source/OpenMldb.md
b/docs/zh/connectors/source/OpenMldb.md
index f36349aac3..1f5eaf20ab 100644
--- a/docs/zh/connectors/source/OpenMldb.md
+++ b/docs/zh/connectors/source/OpenMldb.md
@@ -15,6 +15,9 @@ import ChangeLog from '../changelog/connector-openmldb.md';
用于从 OpenMLDB 读取数据。连接器会执行配置的 SQL 语句并把结果转换为 SeaTunnel 记录,同时支持
单机版和集群版两种部署模式。
+查询直接读取在线表。集群读取不会提交离线 Spark 作业、返回作业元数据或修改 OpenMLDB 会话及全局执行模式。
+此连接器不支持读取离线特征数据。
+
## 关键特性
- [x] [批处理](../../introduction/concepts/connector-v2-features.md)
@@ -23,11 +26,12 @@ import ChangeLog from '../changelog/connector-openmldb.md';
- [x] [列投影](../../introduction/concepts/connector-v2-features.md)
- [ ] [并行度](../../introduction/concepts/connector-v2-features.md)
- [ ] [支持用户自定义分片](../../introduction/concepts/connector-v2-features.md)
+- [x] [支持多表读取](../../introduction/concepts/connector-v2-features.md)
## 数据类型映射
-OpenMLDB 类型会按照所配置 `sql` 语句的结果集映射为 SeaTunnel 类型。SeaTunnel 不原生支持的类型会直接
-导致读取失败,并抛出 `UNSUPPORTED_DATA_TYPE` 错误。
+多表模式下,`schema.fields` 声明每个查询返回的字段名称和类型,连接器在输出数据前校验查询结果。
+SQL `NULL` 值保留为 `null`,包括数值和布尔类型,不会被转换为零或 `false`。
| OpenMLDB 数据类型 | SeaTunnel 数据类型 |
|-------------------|--------------------|
@@ -45,7 +49,8 @@ OpenMLDB 类型会按照所配置 `sql` 语句的结果集映射为 SeaTunnel
| 名称 | 类型 | 必需 | 默认值 | 描述
|
|-----------------|---------|------|--------|---------------------------------------------------------------------------------------------------|
| cluster_mode | boolean | 是 | - | 是否以 OpenMLDB 集群模式连接。`false`
表示单机模式,`true` 表示集群模式。 |
-| sql | string | 是 | - | 用于读取数据的 SQL 语句,列名和类型按结果集定义。
|
+| sql | string | 条件必填 | - | 单查询 SQL。与 `tables_configs` 必须二选一。 |
+| tables_configs | list | 条件必填 | - | 多表读取的查询及显式结果结构。 |
| database | string | 是 | - | 要连接的 OpenMLDB 数据库名称。
|
| host | string | 否 | - | 当 `cluster_mode` 为 `false`
时必填,OpenMLDB 单机版主机地址。 |
| port | int | 否 | - | 当 `cluster_mode` 为 `false`
时必填,OpenMLDB 单机版端口。 |
@@ -62,9 +67,30 @@ OpenMLDB 类型会按照所配置 `sql` 语句的结果集映射为 SeaTunnel
### sql [string]
-无论使用单机模式还是集群模式,必填项 `sql` 都不能为空字符串或仅包含空白字符。
+未配置 `tables_configs` 时,`sql` 必填,且不能为空字符串或仅包含空白字符。
+
+该兼容模式使用 SDK 的输入结构接口获取表结构,结果列的数量、顺序和类型必须与输入结构一致。
+需要列投影或别名时,请使用 `tables_configs` 并显式配置结果结构。
+
+### tables_configs [list]
+
+同一个 OpenMLDB 实例上的非空查询列表。每个配置项包含:
+
+- `sql`:非空 SQL 查询。
+- `database`:可选,用于覆盖必填的根级数据库配置。
+- `schema.table`:唯一的输出表标识,用于下游路由。
+- `schema.fields`:查询返回的全部字段名称及受支持的 SeaTunnel 类型。
+
+字段名称必须与查询结果一致,区分大小写,可通过 SQL 别名进行匹配。
+字段按名称映射,`schema.fields` 的顺序不必与 SQL 投影顺序一致。
+结果列缺失、重复、多余或类型不匹配时,读取失败。
+
+连接和超时选项必须配置在源级别。不得同时配置 `tables_configs` 与根级 `sql` 或 `schema`。
+每个配置项可以使用不同的结果结构。
-针对 OpenMLDB 执行的 SQL 语句,结果集的列会成为连接器输出行的字段。
+单个读取器依次执行所有查询。批模式下,所有查询成功后才报告完成,空结果不会跳过后续查询。
+流模式下,每次轮询都会重新执行全部查询,因此可能产生重复记录;这不是 CDC 或增量轮询。
+多表读取不增加并行度,也不提供跨表一致性快照或精确一次保证。
### database [string]
@@ -155,6 +181,51 @@ sink {
}
```
+### 多表读取
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "BATCH"
+}
+
+source {
+ OpenMldb {
+ cluster_mode = false
+ host = "openmldb"
+ port = 6527
+ database = "shop"
+ tables_configs = [
+ {
+ sql = "select id, amount from orders"
+ schema {
+ table = "shop.orders"
+ fields {
+ id = STRING
+ amount = INT
+ }
+ }
+ },
+ {
+ database = "crm"
+ sql = "select id, name from customers"
+ schema {
+ table = "crm.customers"
+ fields {
+ id = STRING
+ name = STRING
+ }
+ }
+ }
+ ]
+ }
+}
+
+sink {
+ Console {}
+}
+```
+
## 变更日志
<ChangeLog />
diff --git
a/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/config/OpenMldbSqlExecutor.java
b/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/config/OpenMldbSqlExecutor.java
index 0f8154f035..f3787c3ffd 100644
---
a/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/config/OpenMldbSqlExecutor.java
+++
b/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/config/OpenMldbSqlExecutor.java
@@ -17,48 +17,61 @@
package org.apache.seatunnel.connectors.seatunnel.openmldb.config;
+import com._4paradigm.openmldb.SQLRouter;
+import com._4paradigm.openmldb.SQLRouterOptions;
+import com._4paradigm.openmldb.StandaloneOptions;
import com._4paradigm.openmldb.sdk.SdkOption;
import com._4paradigm.openmldb.sdk.SqlException;
import com._4paradigm.openmldb.sdk.impl.SqlClusterExecutor;
+import com._4paradigm.openmldb.sql_router_sdk;
public class OpenMldbSqlExecutor {
- private static final SdkOption SDK_OPTION = new SdkOption();
- private static volatile SqlClusterExecutor SQL_EXECUTOR;
-
private OpenMldbSqlExecutor() {}
- public static void initSdkOption(OpenMldbParameters openMldbParameters) {
- if (openMldbParameters.getClusterMode()) {
- SDK_OPTION.setZkCluster(openMldbParameters.getZkHost());
- SDK_OPTION.setZkPath(openMldbParameters.getZkPath());
- } else {
- SDK_OPTION.setHost(openMldbParameters.getHost());
- SDK_OPTION.setPort(openMldbParameters.getPort());
- SDK_OPTION.setClusterMode(false);
- }
- SDK_OPTION.setSessionTimeout(openMldbParameters.getSessionTimeout());
- SDK_OPTION.setRequestTimeout(openMldbParameters.getRequestTimeout());
+ /** Creates an executor owned by a single schema discovery operation. */
+ public static SqlClusterExecutor create(OpenMldbParameters
openMldbParameters)
+ throws SqlException {
+ return new SqlClusterExecutor(options(openMldbParameters));
}
- public static SqlClusterExecutor getSqlExecutor() throws SqlException {
- if (SQL_EXECUTOR == null) {
- synchronized (OpenMldbSqlExecutor.class) {
- if (SQL_EXECUTOR == null) {
- SQL_EXECUTOR = new SqlClusterExecutor(SDK_OPTION);
- }
+ /** Creates a reader-owned router to access the result set's native null
checks. */
+ public static SQLRouter createReader(OpenMldbParameters parameters) throws
SqlException {
+ SqlClusterExecutor.initJavaSdkLibrary("sql_jsdk");
+ SdkOption option = options(parameters);
+ SQLRouter router;
+ if (option.isClusterMode()) {
+ SQLRouterOptions nativeOptions = option.buildSQLRouterOptions();
+ try {
+ router = sql_router_sdk.NewClusterSQLRouter(nativeOptions);
+ } finally {
+ nativeOptions.delete();
+ }
+ } else {
+ StandaloneOptions nativeOptions = option.buildStandaloneOptions();
+ try {
+ router = sql_router_sdk.NewStandaloneSQLRouter(nativeOptions);
+ } finally {
+ nativeOptions.delete();
}
}
- return SQL_EXECUTOR;
+ if (router == null) {
+ throw new SqlException("Failed to create OpenMldb reader");
+ }
+ return router;
}
- public static void close() {
- if (SQL_EXECUTOR != null) {
- synchronized (OpenMldbParameters.class) {
- if (SQL_EXECUTOR != null) {
- SQL_EXECUTOR.close();
- SQL_EXECUTOR = null;
- }
- }
+ private static SdkOption options(OpenMldbParameters openMldbParameters) {
+ SdkOption sdkOption = new SdkOption();
+ if (openMldbParameters.getClusterMode()) {
+ sdkOption.setZkCluster(openMldbParameters.getZkHost());
+ sdkOption.setZkPath(openMldbParameters.getZkPath());
+ } else {
+ sdkOption.setHost(openMldbParameters.getHost());
+ sdkOption.setPort(openMldbParameters.getPort());
+ sdkOption.setClusterMode(false);
}
+ sdkOption.setSessionTimeout(openMldbParameters.getSessionTimeout());
+ sdkOption.setRequestTimeout(openMldbParameters.getRequestTimeout());
+ return sdkOption;
}
}
diff --git
a/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbReadClient.java
b/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbReadClient.java
new file mode 100644
index 0000000000..4f3d669f3f
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbReadClient.java
@@ -0,0 +1,209 @@
+/*
+ * 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.openmldb.source;
+
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import org.apache.seatunnel.api.table.type.SqlType;
+import
org.apache.seatunnel.connectors.seatunnel.openmldb.config.OpenMldbParameters;
+import
org.apache.seatunnel.connectors.seatunnel.openmldb.config.OpenMldbSqlExecutor;
+
+import com._4paradigm.openmldb.DataType;
+import com._4paradigm.openmldb.Date;
+import com._4paradigm.openmldb.ResultSet;
+import com._4paradigm.openmldb.SQLRouter;
+import com._4paradigm.openmldb.Schema;
+import com._4paradigm.openmldb.Status;
+import com._4paradigm.openmldb.sdk.SqlException;
+
+import java.sql.SQLException;
+import java.sql.Timestamp;
+import java.time.LocalDate;
+import java.util.HashMap;
+import java.util.Map;
+
+/** Owns the native resources needed to read nullable query results with SDK
0.6.3. */
+class OpenMldbReadClient implements AutoCloseable {
+ private SQLRouter router;
+
+ OpenMldbReadClient(OpenMldbParameters parameters) throws SqlException {
+ router = OpenMldbSqlExecutor.createReader(parameters);
+ }
+
+ synchronized Query execute(
+ String database, String sql, SeaTunnelRowType rowType, boolean
matchByName)
+ throws SQLException {
+ if (router == null) {
+ throw new SQLException("OpenMldb reader is closed");
+ }
+ Status status = new Status();
+ ResultSet result = null;
+ try {
+ // Query the online tables directly; ExecuteSQL can submit an
offline job in cluster
+ // mode.
+ result = router.ExecuteSQLParameterized(database, sql, null,
status);
+ if (status.getCode() != 0 || result == null) {
+ throw new SQLException("OpenMldb query failed: " +
status.getMsg());
+ }
+ Query query = new Query(result, rowType, matchByName);
+ result = null; // Ownership transfers only after schema validation
succeeds.
+ return query;
+ } finally {
+ if (result != null) {
+ result.delete();
+ }
+ status.delete();
+ }
+ }
+
+ @Override
+ public synchronized void close() {
+ if (router != null) {
+ router.delete();
+ router = null;
+ }
+ }
+
+ static class Query implements AutoCloseable {
+ private ResultSet result;
+ private final SeaTunnelRowType rowType;
+ private final int[] columnIndexes;
+
+ private Query(ResultSet result, SeaTunnelRowType rowType, boolean
matchByName)
+ throws SQLException {
+ this.result = result;
+ this.rowType = rowType;
+ this.columnIndexes = new int[rowType.getTotalFields()];
+ Schema schema = result.GetSchema();
+ if (schema == null) {
+ throw new SQLException("OpenMldb query did not return a
schema");
+ }
+ try {
+ if (schema.GetColumnCnt() != rowType.getTotalFields()) {
+ throw new SQLException("Query column count does not match
source schema");
+ }
+ Map<String, Integer> columns = new HashMap<>();
+ if (matchByName) {
+ for (int i = 0; i < schema.GetColumnCnt(); i++) {
+ String name = schema.GetColumnName(i);
+ if (columns.put(name, i) != null) {
+ throw new SQLException("Ambiguous query column: "
+ name);
+ }
+ }
+ }
+ for (int i = 0; i < rowType.getTotalFields(); i++) {
+ Integer index =
+ matchByName ? columns.get(rowType.getFieldName(i))
: Integer.valueOf(i);
+ if (index == null) {
+ throw new SQLException(
+ "Query is missing configured field: " +
rowType.getFieldName(i));
+ }
+ columnIndexes[i] = index;
+ if (schema.GetColumnType(columnIndexes[i])
+ !=
nativeType(rowType.getFieldType(i).getSqlType())) {
+ throw new SQLException(
+ "Query column "
+ + (i + 1)
+ + " does not match source schema type "
+ +
rowType.getFieldType(i).getSqlType());
+ }
+ }
+ } finally {
+ schema.delete();
+ }
+ }
+
+ boolean next() {
+ return result.Next();
+ }
+
+ SeaTunnelRow readRow() {
+ Object[] fields = new Object[rowType.getTotalFields()];
+ for (int i = 0; i < fields.length; i++) {
+ // JDBC primitive getters return zero/false for NULL, and
wasNull is unsupported.
+ if (!result.IsNULL(columnIndexes[i])) {
+ fields[i] = readValue(columnIndexes[i],
rowType.getFieldType(i).getSqlType());
+ }
+ }
+ return new SeaTunnelRow(fields);
+ }
+
+ private Object readValue(int index, SqlType type) {
+ switch (type) {
+ case BOOLEAN:
+ return result.GetBoolUnsafe(index);
+ case SMALLINT:
+ return result.GetInt16Unsafe(index);
+ case INT:
+ return result.GetInt32Unsafe(index);
+ case BIGINT:
+ return result.GetInt64Unsafe(index);
+ case FLOAT:
+ return result.GetFloatUnsafe(index);
+ case DOUBLE:
+ return result.GetDoubleUnsafe(index);
+ case STRING:
+ return result.GetStringUnsafe(index);
+ case DATE:
+ Date date = result.GetStructDateUnsafe(index);
+ try {
+ return LocalDate.of(date.getYear(), date.getMonth(),
date.getDay());
+ } finally {
+ date.delete();
+ }
+ case TIMESTAMP:
+ return new
Timestamp(result.GetTimeUnsafe(index)).toLocalDateTime();
+ default:
+ throw new IllegalArgumentException("Unsupported OpenMldb
type: " + type);
+ }
+ }
+
+ @Override
+ public void close() {
+ if (result != null) {
+ result.delete();
+ result = null;
+ }
+ }
+ }
+
+ private static DataType nativeType(SqlType type) throws SQLException {
+ switch (type) {
+ case BOOLEAN:
+ return DataType.kTypeBool;
+ case SMALLINT:
+ return DataType.kTypeInt16;
+ case INT:
+ return DataType.kTypeInt32;
+ case BIGINT:
+ return DataType.kTypeInt64;
+ case FLOAT:
+ return DataType.kTypeFloat;
+ case DOUBLE:
+ return DataType.kTypeDouble;
+ case STRING:
+ return DataType.kTypeString;
+ case DATE:
+ return DataType.kTypeDate;
+ case TIMESTAMP:
+ return DataType.kTypeTimestamp;
+ default:
+ throw new SQLException("Unsupported OpenMldb type: " + type);
+ }
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbSource.java
b/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbSource.java
index 7b038544eb..9acb77b8d3 100644
---
a/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbSource.java
+++
b/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbSource.java
@@ -44,32 +44,48 @@ import com._4paradigm.openmldb.sdk.impl.SqlClusterExecutor;
import java.sql.SQLException;
import java.sql.Types;
+import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
public class OpenMldbSource extends AbstractSingleSplitSource<SeaTunnelRow>
implements SupportColumnProjection {
- private final OpenMldbParameters openMldbParameters;
- private final CatalogTable catalogTable;
+ private final List<OpenMldbParameters> tableParameters;
+ private final List<CatalogTable> catalogTables;
+ private final boolean multiTable;
private JobContext jobContext;
public OpenMldbSource(OpenMldbParameters openMldbParameters) {
- this.openMldbParameters = openMldbParameters;
- OpenMldbSqlExecutor.initSdkOption(openMldbParameters);
+ this.tableParameters = Collections.singletonList(openMldbParameters);
+ this.multiTable = false;
+ SqlClusterExecutor sqlExecutor = null;
try {
- SqlClusterExecutor sqlExecutor =
OpenMldbSqlExecutor.getSqlExecutor();
+ sqlExecutor = OpenMldbSqlExecutor.create(openMldbParameters);
Schema inputSchema =
sqlExecutor.getInputSchema(
openMldbParameters.getDatabase(),
openMldbParameters.getSql());
List<Column> columnList = inputSchema.getColumnList();
- this.catalogTable = convert(columnList);
+ this.catalogTables =
+ Collections.singletonList(
+ convert(columnList,
openMldbParameters.getDatabase()));
} catch (SQLException | SqlException e) {
throw new OpenMldbConnectorException(
CommonErrorCodeDeprecated.TABLE_SCHEMA_GET_FAILED,
- "Failed to initialize data schema");
+ "Failed to initialize data schema",
+ e);
+ } finally {
+ if (sqlExecutor != null) {
+ sqlExecutor.close();
+ }
}
}
+ OpenMldbSource(List<OpenMldbParameters> tableParameters,
List<CatalogTable> catalogTables) {
+ this.tableParameters = Collections.unmodifiableList(new
ArrayList<>(tableParameters));
+ this.catalogTables = Collections.unmodifiableList(new
ArrayList<>(catalogTables));
+ this.multiTable = true;
+ }
+
@Override
public String getPluginName() {
return "OpenMldb";
@@ -84,14 +100,13 @@ public class OpenMldbSource extends
AbstractSingleSplitSource<SeaTunnelRow>
@Override
public List<CatalogTable> getProducedCatalogTables() {
- return Collections.singletonList(catalogTable);
+ return catalogTables;
}
@Override
public AbstractSingleSplitReader<SeaTunnelRow> createReader(
SingleSplitReaderContext readerContext) throws Exception {
- return new OpenMldbSourceReader(
- openMldbParameters, catalogTable.getSeaTunnelRowType(),
readerContext);
+ return new OpenMldbSourceReader(tableParameters, catalogTables,
multiTable, readerContext);
}
@Override
@@ -126,7 +141,7 @@ public class OpenMldbSource extends
AbstractSingleSplitSource<SeaTunnelRow>
}
}
- private CatalogTable convert(List<Column> columnList) {
+ private CatalogTable convert(List<Column> columnList, String database) {
TableSchema.Builder builder = TableSchema.builder();
for (int i = 0; i < columnList.size(); i++) {
Column column = columnList.get(i);
@@ -135,15 +150,15 @@ public class OpenMldbSource extends
AbstractSingleSplitSource<SeaTunnelRow>
column.getColumnName(),
convertSeaTunnelDataType(column.getSqlType()),
(Long) null,
- column.isNotNull(),
+ !column.isNotNull(),
null,
null));
}
return CatalogTable.of(
- TableIdentifier.of("OpenMldb",
openMldbParameters.getDatabase(), "default"),
+ TableIdentifier.of("OpenMldb", database, "default"),
builder.build(),
- null,
- null,
+ Collections.emptyMap(),
+ Collections.emptyList(),
null);
}
}
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 0ff80bbac2..a45dbf39e0 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
@@ -17,19 +17,34 @@
package org.apache.seatunnel.connectors.seatunnel.openmldb.source;
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConditionExtension;
+import org.apache.seatunnel.api.configuration.util.Conditions;
+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.options.ConnectorCommonOptions;
import org.apache.seatunnel.api.source.SeaTunnelSource;
import org.apache.seatunnel.api.source.SourceSplit;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
import org.apache.seatunnel.api.table.connector.TableSource;
import org.apache.seatunnel.api.table.factory.Factory;
import org.apache.seatunnel.api.table.factory.TableSourceFactory;
import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
+import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
import
org.apache.seatunnel.connectors.seatunnel.openmldb.config.OpenMldbParameters;
import
org.apache.seatunnel.connectors.seatunnel.openmldb.config.OpenMldbSourceOptions;
import com.google.auto.service.AutoService;
import java.io.Serializable;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
import static org.apache.seatunnel.api.configuration.util.Conditions.notBlank;
@@ -44,7 +59,13 @@ public class OpenMldbSourceFactory implements
TableSourceFactory {
public OptionRule optionRule() {
return OptionRule.builder()
.required(OpenMldbSourceOptions.CLUSTER_MODE)
- .required(OpenMldbSourceOptions.SQL,
notBlank(OpenMldbSourceOptions.SQL))
+ .exclusive(OpenMldbSourceOptions.SQL,
ConnectorCommonOptions.TABLE_CONFIGS)
+ .optional(OpenMldbSourceOptions.SQL,
notBlank(OpenMldbSourceOptions.SQL))
+ .optional(
+ ConnectorCommonOptions.TABLE_CONFIGS,
+
Conditions.notEmpty(ConnectorCommonOptions.TABLE_CONFIGS),
+ Conditions.extension(
+ ConnectorCommonOptions.TABLE_CONFIGS, new
TablesValidator()))
.required(OpenMldbSourceOptions.DATABASE)
.optional(OpenMldbSourceOptions.SESSION_TIMEOUT)
.optional(OpenMldbSourceOptions.REQUEST_TIMEOUT)
@@ -69,8 +90,110 @@ public class OpenMldbSourceFactory implements
TableSourceFactory {
@Override
public <T, SplitT extends SourceSplit, StateT extends Serializable>
TableSource<T, SplitT, StateT>
createSource(TableSourceFactoryContext context) {
+ ConfigValidator.of(context.getOptions()).validate(optionRule());
+ if
(context.getOptions().getOptional(ConnectorCommonOptions.TABLE_CONFIGS).isPresent())
{
+ List<OpenMldbParameters> tables =
buildTables(context.getOptions());
+ List<CatalogTable> catalogTables = new ArrayList<>();
+ for (Map<String, Object> entry :
+
context.getOptions().get(ConnectorCommonOptions.TABLE_CONFIGS)) {
+ catalogTables.add(
+ CatalogTableUtil.buildWithConfig(
+ "OpenMldb", ReadonlyConfig.fromMap(entry)));
+ }
+ return () ->
+ (SeaTunnelSource<T, SplitT, StateT>) new
OpenMldbSource(tables, catalogTables);
+ }
OpenMldbParameters openMldbParameters =
OpenMldbParameters.buildWithConfig(context.getOptions().toConfig());
return () -> (SeaTunnelSource<T, SplitT, StateT>) new
OpenMldbSource(openMldbParameters);
}
+
+ private static List<OpenMldbParameters> buildTables(ReadonlyConfig config)
{
+ List<OpenMldbParameters> tables = new ArrayList<>();
+ for (Map<String, Object> entry :
config.get(ConnectorCommonOptions.TABLE_CONFIGS)) {
+ tables.add(
+ OpenMldbParameters.buildWithConfig(
+ ReadonlyConfig.fromMap(entry)
+ .toConfig()
+ .withFallback(
+ config.toConfig()
+ .withoutPath(
+
ConnectorCommonOptions.TABLE_CONFIGS
+ .key()))));
+ }
+ return tables;
+ }
+
+ static class TablesValidator implements
ConditionExtension<List<Map<String, Object>>> {
+ @Override
+ public String description() {
+ return "each table requires non-blank sql and an explicit schema
with a unique table identity";
+ }
+
+ @Override
+ public boolean evaluate(ReadonlyConfig config, List<Map<String,
Object>> entries) {
+ if (config.getOptional(ConnectorCommonOptions.SCHEMA).isPresent())
{
+ throw new OptionValidationException(
+ "With tables_configs, configure schema.table inside
each entry");
+ }
+ Set<String> ids = new HashSet<>();
+ Set<String> supported = new HashSet<>(Arrays.asList("sql",
"database", "schema"));
+ for (int i = 0; i < entries.size(); i++) {
+ Map<String, Object> entry = entries.get(i);
+ if (entry == null || !supported.containsAll(entry.keySet())) {
+ throw new OptionValidationException(
+ "tables_configs[%d]: only sql, database and schema
are supported; connection options belong at source level",
+ i);
+ }
+ Object schema = entry.get("schema");
+ if (!(schema instanceof Map)
+ || !(((Map<?, ?>) schema).get("fields") instanceof Map)
+ || ((Map<?, ?>) ((Map<?, ?>)
schema).get("fields")).isEmpty()) {
+ throw new OptionValidationException(
+ "tables_configs[%d]: schema.table and non-empty
schema.fields are required",
+ i);
+ }
+ Object table = ((Map<?, ?>) schema).get("table");
+ Object database =
+ entry.containsKey("database")
+ ? entry.get("database")
+ : config.get(OpenMldbSourceOptions.DATABASE);
+ if (!nonBlank(entry.get("sql")) || !nonBlank(table) ||
!nonBlank(database)) {
+ throw new OptionValidationException(
+ "tables_configs[%d]: sql, database and
schema.table must be non-blank strings",
+ i);
+ }
+ CatalogTable catalogTable =
+ CatalogTableUtil.buildWithConfig("OpenMldb",
ReadonlyConfig.fromMap(entry));
+ for (SeaTunnelDataType<?> type :
+ catalogTable.getSeaTunnelRowType().getFieldTypes()) {
+ switch (type.getSqlType()) {
+ case BOOLEAN:
+ case SMALLINT:
+ case INT:
+ case BIGINT:
+ case FLOAT:
+ case DOUBLE:
+ case STRING:
+ case DATE:
+ case TIMESTAMP:
+ break;
+ default:
+ throw new OptionValidationException(
+ "tables_configs[%d]: unsupported OpenMldb
type '%s'", i, type);
+ }
+ }
+ String id = catalogTable.getTableId().toTablePath().toString();
+ if (!ids.add(id)) {
+ throw new OptionValidationException(
+ "tables_configs[%d]: duplicate table identity
'%s'", i, id);
+ }
+ }
+ return true;
+ }
+
+ private boolean nonBlank(Object value) {
+ return value instanceof String && !((String)
value).trim().isEmpty();
+ }
+ }
}
diff --git
a/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbSourceReader.java
b/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbSourceReader.java
index c034fd68eb..9082756b52 100644
---
a/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbSourceReader.java
+++
b/seatunnel-connectors-v2/connector-openmldb/src/main/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbSourceReader.java
@@ -19,101 +19,108 @@ package
org.apache.seatunnel.connectors.seatunnel.openmldb.source;
import org.apache.seatunnel.api.source.Boundedness;
import org.apache.seatunnel.api.source.Collector;
-import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
import org.apache.seatunnel.api.table.type.SeaTunnelRow;
import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
import org.apache.seatunnel.common.exception.CommonErrorCodeDeprecated;
import
org.apache.seatunnel.connectors.seatunnel.common.source.AbstractSingleSplitReader;
import
org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplitReaderContext;
import
org.apache.seatunnel.connectors.seatunnel.openmldb.config.OpenMldbParameters;
-import
org.apache.seatunnel.connectors.seatunnel.openmldb.config.OpenMldbSqlExecutor;
import
org.apache.seatunnel.connectors.seatunnel.openmldb.exception.OpenMldbConnectorException;
-import com._4paradigm.openmldb.sdk.impl.SqlClusterExecutor;
import lombok.extern.slf4j.Slf4j;
import java.io.IOException;
-import java.sql.Date;
-import java.sql.ResultSet;
import java.sql.SQLException;
-import java.sql.Timestamp;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
@Slf4j
public class OpenMldbSourceReader extends
AbstractSingleSplitReader<SeaTunnelRow> {
- private final OpenMldbParameters openMldbParameters;
- private final SeaTunnelRowType seaTunnelRowType;
+ private final List<OpenMldbParameters> tableParameters;
+ private final List<SeaTunnelRowType> rowTypes;
+ private final List<String> tableIds;
private final SingleSplitReaderContext readerContext;
+ private OpenMldbReadClient client;
+ private boolean finished;
public OpenMldbSourceReader(
OpenMldbParameters openMldbParameters,
SeaTunnelRowType seaTunnelRowType,
SingleSplitReaderContext readerContext) {
- this.openMldbParameters = openMldbParameters;
- this.seaTunnelRowType = seaTunnelRowType;
+ this.tableParameters = Collections.singletonList(openMldbParameters);
+ this.rowTypes = Collections.singletonList(seaTunnelRowType);
+ this.tableIds = Collections.singletonList(null);
+ this.readerContext = readerContext;
+ }
+
+ OpenMldbSourceReader(
+ List<OpenMldbParameters> tableParameters,
+ List<CatalogTable> catalogTables,
+ boolean multiTable,
+ SingleSplitReaderContext readerContext) {
+ this.tableParameters = new ArrayList<>(tableParameters);
+ this.rowTypes = new ArrayList<>();
+ this.tableIds = new ArrayList<>();
+ for (CatalogTable table : catalogTables) {
+ rowTypes.add(table.getSeaTunnelRowType());
+ tableIds.add(multiTable ?
table.getTableId().toTablePath().toString() : null);
+ }
this.readerContext = readerContext;
}
@Override
public void open() throws Exception {
- OpenMldbSqlExecutor.initSdkOption(openMldbParameters);
+ client = new OpenMldbReadClient(tableParameters.get(0));
}
@Override
public void close() throws IOException {
- OpenMldbSqlExecutor.close();
+ if (client != null) {
+ client.close();
+ client = null;
+ }
}
@Override
public void pollNext(Collector<SeaTunnelRow> output) throws Exception {
- int totalFields = seaTunnelRowType.getTotalFields();
- Object[] objects = new Object[totalFields];
- SqlClusterExecutor sqlExecutor = OpenMldbSqlExecutor.getSqlExecutor();
- try (ResultSet resultSet =
- sqlExecutor.executeSQL(
- openMldbParameters.getDatabase(),
openMldbParameters.getSql())) {
- while (resultSet.next()) {
- for (int i = 0; i < totalFields; i++) {
- objects[i] = getObject(resultSet, i,
seaTunnelRowType.getFieldType(i));
- }
- output.collect(new SeaTunnelRow(objects));
- }
- } finally {
- if (Boundedness.BOUNDED.equals(readerContext.getBoundedness())) {
- // signal to the source that we have reached the end of the
data.
- log.info("Closed the bounded openmldb source");
- readerContext.signalNoMoreElement();
- }
+ if (finished) {
+ return;
+ }
+ for (int i = 0; i < tableParameters.size(); i++) {
+ readTable(output, tableParameters.get(i), rowTypes.get(i),
tableIds.get(i));
+ }
+ if (Boundedness.BOUNDED.equals(readerContext.getBoundedness())) {
+ finished = true;
+ log.info("Finished reading the bounded OpenMldb source");
+ readerContext.signalNoMoreElement();
}
}
- private Object getObject(ResultSet resultSet, int index,
SeaTunnelDataType<?> dataType)
+ private void readTable(
+ Collector<SeaTunnelRow> output,
+ OpenMldbParameters parameters,
+ SeaTunnelRowType rowType,
+ String tableId)
throws SQLException {
- index = index + 1;
- switch (dataType.getSqlType()) {
- case BOOLEAN:
- return resultSet.getBoolean(index);
- case INT:
- return resultSet.getInt(index);
- case SMALLINT:
- return resultSet.getShort(index);
- case BIGINT:
- return resultSet.getLong(index);
- case FLOAT:
- return resultSet.getFloat(index);
- case DOUBLE:
- return resultSet.getDouble(index);
- case STRING:
- return resultSet.getString(index);
- case DATE:
- Date date = resultSet.getDate(index);
- return date.toLocalDate();
- case TIMESTAMP:
- Timestamp timestamp = resultSet.getTimestamp(index);
- return timestamp.toLocalDateTime();
- default:
- throw new OpenMldbConnectorException(
- CommonErrorCodeDeprecated.UNSUPPORTED_DATA_TYPE,
- "Unsupported this data type");
+ try (OpenMldbReadClient.Query query =
+ client.execute(
+ parameters.getDatabase(), parameters.getSql(),
rowType, tableId != null)) {
+ while (query.next()) {
+ SeaTunnelRow row = query.readRow();
+ if (tableId != null) {
+ row.setTableId(tableId);
+ }
+ output.collect(row);
+ }
+ } catch (SQLException | RuntimeException e) {
+ throw new OpenMldbConnectorException(
+ CommonErrorCodeDeprecated.READER_OPERATION_FAILED,
+ "Failed to read OpenMldb table '"
+ + (tableId == null ? parameters.getDatabase() :
tableId)
+ + "'",
+ e);
}
}
}
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 daccae7527..a1a51954d7 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
@@ -28,6 +28,8 @@ import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource;
+import java.util.Arrays;
+import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
@@ -35,6 +37,137 @@ class OpenMldbFactoryTest {
private final OptionRule optionRule = new
OpenMldbSourceFactory().optionRule();
+ @Test
+ void testMultiTableConfigAccepted() {
+ Map<String, Object> config = requiredConfig(false);
+ config.remove("sql");
+ config.put(
+ "tables_configs",
+ Arrays.asList(
+ table("orders", "select * from orders"),
+ table("customers", "select * from customers")));
+ Assertions.assertDoesNotThrow(() -> validate(config));
+ }
+
+ private Map<String, Object> table(String name, String sql) {
+ Map<String, Object> table = new HashMap<>();
+ table.put("sql", sql);
+ Map<String, Object> schema = new HashMap<>();
+ schema.put("table", name);
+ schema.put("fields", Collections.singletonMap("id", "STRING"));
+ table.put("schema", schema);
+ return table;
+ }
+
+ @Test
+ void testSqlAndTablesAreMutuallyExclusive() {
+ Map<String, Object> config = requiredConfig(false);
+ config.put(
+ "tables_configs",
+ Collections.singletonList(table("orders", "select * from
orders")));
+ Assertions.assertThrows(OptionValidationException.class, () ->
validate(config));
+ }
+
+ @Test
+ void testEmptyTablesRejected() {
+ Map<String, Object> config = multiConfig();
+ config.put("tables_configs", Collections.emptyList());
+ Assertions.assertThrows(OptionValidationException.class, () ->
validate(config));
+ }
+
+ @Test
+ void testDuplicateTableIdentityRejected() {
+ Map<String, Object> config = multiConfig();
+ config.put(
+ "tables_configs",
+ Arrays.asList(
+ table("orders", "select * from orders"),
+ table("orders", "select * from archived")));
+ Assertions.assertThrows(OptionValidationException.class, () ->
validate(config));
+ }
+
+ @Test
+ void testEntrySqlRequiredAndNotBlank() {
+ for (String sql : new String[] {null, "", " \t\n"}) {
+ Map<String, Object> config = multiConfig();
+ config.put("tables_configs",
Collections.singletonList(table("orders", sql)));
+ Assertions.assertThrows(OptionValidationException.class, () ->
validate(config));
+ }
+ }
+
+ @Test
+ void testSchemaFieldsRequired() {
+ for (Object fields : new Object[] {null, Collections.emptyMap(), "id
STRING"}) {
+ Map<String, Object> config = multiConfig();
+ Map<String, Object> entry = table("orders", "select * from
orders");
+ ((Map<String, Object>) entry.get("schema")).put("fields", fields);
+ config.put("tables_configs", Collections.singletonList(entry));
+ Assertions.assertThrows(OptionValidationException.class, () ->
validate(config));
+ }
+ }
+
+ @Test
+ void testTableIdentityRequired() {
+ for (String tableId : new String[] {null, "", " \t"}) {
+ Map<String, Object> config = multiConfig();
+ config.put(
+ "tables_configs",
+ Collections.singletonList(table(tableId, "select * from
orders")));
+ Assertions.assertThrows(OptionValidationException.class, () ->
validate(config));
+ }
+ }
+
+ @Test
+ void testRootSchemaRejectedInMultiTableMode() {
+ Map<String, Object> config = multiConfig();
+ config.put("schema", table("root", "select * from
orders").get("schema"));
+ Assertions.assertThrows(OptionValidationException.class, () ->
validate(config));
+ }
+
+ @Test
+ void testEntryConnectionOverridesRejected() {
+ for (String key :
+ new String[] {"host", "port", "cluster_mode", "zk_host",
"request_timeout"}) {
+ Map<String, Object> config = multiConfig();
+ Map<String, Object> entry = table("orders", "select * from
orders");
+ entry.put(key, "override");
+ config.put("tables_configs", Collections.singletonList(entry));
+ Assertions.assertThrows(OptionValidationException.class, () ->
validate(config));
+ }
+ }
+
+ @Test
+ void testEntryDatabaseOverrideValidated() {
+ Map<String, Object> config = multiConfig();
+ Map<String, Object> entry = table("orders", "select * from orders");
+ config.put("tables_configs", Collections.singletonList(entry));
+ entry.put("database", "another_db");
+ Assertions.assertDoesNotThrow(() -> validate(config));
+ entry.put("database", " ");
+ Assertions.assertThrows(OptionValidationException.class, () ->
validate(config));
+ }
+
+ private Map<String, Object> multiConfig() {
+ Map<String, Object> config = requiredConfig(false);
+ config.remove("sql");
+ config.put(
+ "tables_configs",
+ Collections.singletonList(table("orders", "select * from
orders")));
+ return config;
+ }
+
+ @Test
+ void testUnsupportedFieldTypesRejectedBeforeConnecting() {
+ for (String type : new String[] {"DECIMAL(10,2)", "TINYINT",
"ARRAY<INT>", "BYTES"}) {
+ Map<String, Object> config = multiConfig();
+ Map<String, Object> entry = table("orders", "select * from
orders");
+ ((Map<String, Object>) entry.get("schema"))
+ .put("fields", Collections.singletonMap("id", type));
+ config.put("tables_configs", Collections.singletonList(entry));
+ Assertions.assertThrows(OptionValidationException.class, () ->
validate(config));
+ }
+ }
+
@ParameterizedTest
@ValueSource(booleans = {false, true})
void testNonblankSqlAccepted(boolean clusterMode) {
diff --git
a/seatunnel-connectors-v2/connector-openmldb/src/test/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbSourceReaderTest.java
b/seatunnel-connectors-v2/connector-openmldb/src/test/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbSourceReaderTest.java
new file mode 100644
index 0000000000..8cad021ff7
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-openmldb/src/test/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbSourceReaderTest.java
@@ -0,0 +1,238 @@
+/*
+ * 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.openmldb.source;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.source.Boundedness;
+import org.apache.seatunnel.api.source.Collector;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import
org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplitReaderContext;
+import
org.apache.seatunnel.connectors.seatunnel.openmldb.config.OpenMldbParameters;
+import
org.apache.seatunnel.connectors.seatunnel.openmldb.exception.OpenMldbConnectorException;
+
+import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+import org.mockito.MockedConstruction;
+
+import java.sql.SQLException;
+import java.util.Arrays;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyBoolean;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
+import static org.mockito.Mockito.when;
+
+class OpenMldbSourceReaderTest {
+ private final CatalogTable first = table("first", "id = STRING");
+ private final CatalogTable second = table("second", "value = INT");
+ private final SingleSplitReaderContext context =
mock(SingleSplitReaderContext.class);
+ private final Collector<SeaTunnelRow> output = mock(Collector.class);
+
+ @Test
+ void readsDifferentSchemasAndCompletesOnlyOnce() throws Exception {
+ OpenMldbReadClient.Query firstQuery = query(new SeaTunnelRow(new
Object[] {"key"}));
+ OpenMldbReadClient.Query secondQuery = query(new SeaTunnelRow(new
Object[] {null}));
+ when(context.getBoundedness()).thenReturn(Boundedness.BOUNDED);
+ try (MockedConstruction<OpenMldbReadClient> clients =
+ mockConstruction(
+ OpenMldbReadClient.class,
+ (client, ignored) -> {
+ when(client.execute(
+ "one",
+ "select * from first",
+ first.getSeaTunnelRowType(),
+ true))
+ .thenReturn(firstQuery);
+ when(client.execute(
+ "two",
+ "select * from second",
+ second.getSeaTunnelRowType(),
+ true))
+ .thenReturn(secondQuery);
+ })) {
+ OpenMldbSourceReader reader = reader();
+ reader.open();
+ reader.pollNext(output);
+ reader.pollNext(output);
+ ArgumentCaptor<SeaTunnelRow> rows =
ArgumentCaptor.forClass(SeaTunnelRow.class);
+ verify(output, times(2)).collect(rows.capture());
+ assertEquals(
+ first.getTableId().toTablePath().toString(),
+ rows.getAllValues().get(0).getTableId());
+ assertEquals(
+ second.getTableId().toTablePath().toString(),
+ rows.getAllValues().get(1).getTableId());
+ assertNull(rows.getAllValues().get(1).getField(0));
+ verify(context).signalNoMoreElement();
+ verify(firstQuery).close();
+ verify(secondQuery).close();
+ reader.close();
+ reader.close();
+ verify(clients.constructed().get(0)).close();
+ }
+ }
+
+ @Test
+ void emptyFirstTableDoesNotSkipSecondTable() throws Exception {
+ OpenMldbReadClient.Query empty = mock(OpenMldbReadClient.Query.class);
+ OpenMldbReadClient.Query populated = query(new SeaTunnelRow(new
Object[] {7}));
+ try (MockedConstruction<OpenMldbReadClient> ignored =
+ mockConstruction(
+ OpenMldbReadClient.class,
+ (client, construction) -> {
+ when(client.execute(anyString(), anyString(),
any(), anyBoolean()))
+ .thenReturn(empty, populated);
+ })) {
+ try (OpenMldbSourceReader reader = reader()) {
+ reader.open();
+ reader.pollNext(output);
+ }
+ verify(output).collect(any(SeaTunnelRow.class));
+ verify(empty).close();
+ verify(populated).close();
+ }
+ }
+
+ @Test
+ void queryFailureDoesNotReportSuccessfulCompletion() throws Exception {
+ when(context.getBoundedness()).thenReturn(Boundedness.BOUNDED);
+ try (MockedConstruction<OpenMldbReadClient> ignored =
+ mockConstruction(
+ OpenMldbReadClient.class,
+ (client, construction) -> {
+ when(client.execute(anyString(), anyString(),
any(), anyBoolean()))
+ .thenThrow(new SQLException("query
rejected"));
+ })) {
+ try (OpenMldbSourceReader reader = reader()) {
+ reader.open();
+ OpenMldbConnectorException error =
+ assertThrows(
+ OpenMldbConnectorException.class, () ->
reader.pollNext(output));
+ assertTrue(error.getMessage().contains("first"));
+ assertInstanceOf(SQLException.class, error.getCause());
+ }
+ verify(context, never()).signalNoMoreElement();
+ verifyNoInteractions(output);
+ }
+ }
+
+ @Test
+ void collectorFailureClosesQueryAndDoesNotComplete() throws Exception {
+ OpenMldbReadClient.Query query = query(new SeaTunnelRow(new Object[]
{"key"}));
+ doThrow(new IllegalStateException("collector failed"))
+ .when(output)
+ .collect(any(SeaTunnelRow.class));
+ try (MockedConstruction<OpenMldbReadClient> ignored =
+ mockConstruction(
+ OpenMldbReadClient.class,
+ (client, construction) -> {
+ when(client.execute(anyString(), anyString(),
any(), anyBoolean()))
+ .thenReturn(query);
+ })) {
+ try (OpenMldbSourceReader reader = reader()) {
+ reader.open();
+ assertThrows(OpenMldbConnectorException.class, () ->
reader.pollNext(output));
+ }
+ verify(query).close();
+ verify(context, never()).signalNoMoreElement();
+ }
+ }
+
+ @Test
+ void streamingRepeatsQueriesWithoutCompleting() throws Exception {
+ when(context.getBoundedness()).thenReturn(Boundedness.UNBOUNDED);
+ try (MockedConstruction<OpenMldbReadClient> clients =
+ mockConstruction(
+ OpenMldbReadClient.class,
+ (client, construction) -> {
+ when(client.execute(anyString(), anyString(),
any(), anyBoolean()))
+ .thenAnswer(invocation ->
mock(OpenMldbReadClient.Query.class));
+ })) {
+ try (OpenMldbSourceReader reader = reader()) {
+ reader.open();
+ reader.pollNext(output);
+ reader.pollNext(output);
+ verify(clients.constructed().get(0), times(4))
+ .execute(anyString(), anyString(), any(),
anyBoolean());
+ }
+ verify(context, never()).signalNoMoreElement();
+ }
+ }
+
+ @Test
+ void closingOneReaderDoesNotCloseAnother() throws Exception {
+ try (MockedConstruction<OpenMldbReadClient> clients =
+ mockConstruction(OpenMldbReadClient.class)) {
+ OpenMldbSourceReader firstReader = reader();
+ OpenMldbSourceReader secondReader = reader();
+ firstReader.open();
+ secondReader.open();
+ firstReader.close();
+ verify(clients.constructed().get(0)).close();
+ verifyNoInteractions(clients.constructed().get(1));
+ secondReader.close();
+ verify(clients.constructed().get(1)).close();
+ }
+ }
+
+ private OpenMldbSourceReader reader() {
+ return new OpenMldbSourceReader(
+ Arrays.asList(parameters("one", "first"), parameters("two",
"second")),
+ Arrays.asList(first, second),
+ true,
+ context);
+ }
+
+ private static OpenMldbReadClient.Query query(SeaTunnelRow row) {
+ OpenMldbReadClient.Query query = mock(OpenMldbReadClient.Query.class);
+ when(query.next()).thenReturn(true, false);
+ when(query.readRow()).thenReturn(row);
+ return query;
+ }
+
+ private static CatalogTable table(String name, String fields) {
+ return CatalogTableUtil.buildWithConfig(
+ ReadonlyConfig.fromConfig(
+
org.apache.seatunnel.shade.com.typesafe.config.ConfigFactory.parseString(
+ "schema { table = " + name + ", fields { " +
fields + " } }")));
+ }
+
+ private static OpenMldbParameters parameters(String database, String
table) {
+ return OpenMldbParameters.buildWithConfig(
+
org.apache.seatunnel.shade.com.typesafe.config.ConfigFactory.parseString(
+ "cluster_mode = false\nhost = localhost\nport =
6527\ndatabase = "
+ + database
+ + "\nsql = \"select * from "
+ + table
+ + "\""));
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-openmldb/src/test/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbSourceTest.java
b/seatunnel-connectors-v2/connector-openmldb/src/test/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbSourceTest.java
new file mode 100644
index 0000000000..80509d2fac
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-openmldb/src/test/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/OpenMldbSourceTest.java
@@ -0,0 +1,143 @@
+/*
+ * 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.openmldb.source;
+
+import org.apache.seatunnel.shade.com.typesafe.config.ConfigFactory;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplit;
+import
org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplitEnumeratorState;
+import
org.apache.seatunnel.connectors.seatunnel.openmldb.config.OpenMldbParameters;
+import
org.apache.seatunnel.connectors.seatunnel.openmldb.exception.OpenMldbConnectorException;
+
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedConstruction;
+
+import com._4paradigm.openmldb.sdk.Column;
+import com._4paradigm.openmldb.sdk.Schema;
+import com._4paradigm.openmldb.sdk.impl.SqlClusterExecutor;
+
+import java.io.ByteArrayOutputStream;
+import java.io.ObjectOutputStream;
+import java.sql.SQLException;
+import java.sql.Types;
+import java.util.Arrays;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+class OpenMldbSourceTest {
+ @Test
+ void discoveryPreservesNullabilityAndClosesExecutor() throws Exception {
+ Schema schema =
+ new Schema(
+ Arrays.asList(
+ new Column("id", Types.VARCHAR, true, false),
+ new Column("value", Types.INTEGER, false,
false)));
+ try (MockedConstruction<SqlClusterExecutor> executors =
+ mockConstruction(
+ SqlClusterExecutor.class,
+ (executor, context) ->
+ when(executor.getInputSchema(anyString(),
anyString()))
+ .thenReturn(schema))) {
+ OpenMldbSource source = new OpenMldbSource(parameters());
+ assertEquals(1, source.getProducedCatalogTables().size());
+ assertFalse(
+ source.getProducedCatalogTables()
+ .get(0)
+ .getTableSchema()
+ .getColumns()
+ .get(0)
+ .isNullable());
+ assertTrue(
+ source.getProducedCatalogTables()
+ .get(0)
+ .getTableSchema()
+ .getColumns()
+ .get(1)
+ .isNullable());
+ verify(executors.constructed().get(0)).close();
+ }
+ }
+
+ @Test
+ void failedDiscoveryClosesExecutorAndPreservesCause() throws Exception {
+ SQLException cause = new SQLException("invalid query");
+ try (MockedConstruction<SqlClusterExecutor> executors =
+ mockConstruction(
+ SqlClusterExecutor.class,
+ (executor, context) ->
+ when(executor.getInputSchema(anyString(),
anyString()))
+ .thenThrow(cause))) {
+ OpenMldbConnectorException error =
+ assertThrows(
+ OpenMldbConnectorException.class,
+ () -> new OpenMldbSource(parameters()));
+ assertEquals(cause, error.getCause());
+ verify(executors.constructed().get(0)).close();
+ }
+ }
+
+ @Test
+ void multiTableMetadataDoesNotConnectAndSourceIsSerializable() throws
Exception {
+ ReadonlyConfig config =
+ ReadonlyConfig.fromConfig(
+ ConfigFactory.parseString(
+
"cluster_mode=false\nhost=unused\nport=6527\ndatabase=test\n"
+ + "tables_configs=["
+ + "{sql=\"select id from
orders\",schema{table=orders,fields{id=STRING}}},"
+ + "{sql=\"select amount from
sales\",database=other,schema{table=sales,fields{amount=INT}}}"
+ + "]"));
+ try (MockedConstruction<SqlClusterExecutor> executors =
+ mockConstruction(SqlClusterExecutor.class)) {
+ OpenMldbSource source =
+ (OpenMldbSource)
+ new OpenMldbSourceFactory()
+ .<SeaTunnelRow, SingleSplit,
SingleSplitEnumeratorState>
+ createSource(
+ new
TableSourceFactoryContext(
+ config,
getClass().getClassLoader()))
+ .createSource();
+ assertEquals(2, source.getProducedCatalogTables().size());
+ assertEquals(
+ "orders",
+
source.getProducedCatalogTables().get(0).getTableId().toTablePath().toString());
+ assertEquals(
+ "sales",
+
source.getProducedCatalogTables().get(1).getTableId().toTablePath().toString());
+ assertTrue(executors.constructed().isEmpty());
+ try (ObjectOutputStream out = new ObjectOutputStream(new
ByteArrayOutputStream())) {
+ out.writeObject(source);
+ }
+ }
+ }
+
+ private OpenMldbParameters parameters() {
+ return OpenMldbParameters.buildWithConfig(
+ ConfigFactory.parseString(
+
"cluster_mode=false\nhost=localhost\nport=6527\ndatabase=test\nsql=\"select *
from values_table\""));
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-openmldb/src/test/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/TestOpenMldbSourceIT.java
b/seatunnel-connectors-v2/connector-openmldb/src/test/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/TestOpenMldbSourceIT.java
new file mode 100644
index 0000000000..98c2b52087
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-openmldb/src/test/java/org/apache/seatunnel/connectors/seatunnel/openmldb/source/TestOpenMldbSourceIT.java
@@ -0,0 +1,386 @@
+/*
+ * 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.openmldb.source;
+
+import org.apache.seatunnel.shade.com.typesafe.config.ConfigFactory;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.source.Boundedness;
+import org.apache.seatunnel.api.source.Collector;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplit;
+import
org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplitEnumeratorState;
+import
org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplitReaderContext;
+import
org.apache.seatunnel.connectors.seatunnel.openmldb.config.OpenMldbParameters;
+import
org.apache.seatunnel.connectors.seatunnel.openmldb.config.OpenMldbSqlExecutor;
+import
org.apache.seatunnel.connectors.seatunnel.openmldb.exception.OpenMldbConnectorException;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.condition.EnabledIfSystemProperty;
+
+import com._4paradigm.openmldb.sdk.impl.SqlClusterExecutor;
+
+import java.sql.SQLException;
+import java.sql.Timestamp;
+import java.time.LocalDate;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.UUID;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotSame;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * Runs against a disposable OpenMLDB 0.6.3 server. Enable with
-Dopenmldb.integration=true and
+ * optionally -Dopenmldb.host / -Dopenmldb.port for standalone mode, or
-Dopenmldb.cluster=true
+ * -Dopenmldb.zk.host -Dopenmldb.zk.path for cluster mode. The SDK's native
library requires a
+ * compatible Linux amd64 runtime.
+ */
+@EnabledIfSystemProperty(named = "openmldb.integration", matches = "true")
+class TestOpenMldbSourceIT {
+ private static final String DATABASE = "st_" +
UUID.randomUUID().toString().replace("-", "");
+ private static final String SECOND_DATABASE = DATABASE + "_second";
+ private static final String FIELDS =
+ "id=STRING, b=BOOLEAN, s=SMALLINT, i=INT, l=BIGINT, "
+ + "f=FLOAT, d=DOUBLE, text=STRING, day=DATE, ts=TIMESTAMP";
+ private static SqlClusterExecutor setup;
+
+ @BeforeAll
+ static void prepare() throws Exception {
+ setup = OpenMldbSqlExecutor.create(parameters("select * from
values_table"));
+ assertTrue(setup.createDB(DATABASE));
+ assertTrue(setup.createDB(SECOND_DATABASE));
+ String columns =
+ "(id string, b bool, s smallint, i int, l bigint, f float, d
double, "
+ + "text string, day date, ts timestamp,
index(key=id));";
+ assertTrue(setup.executeDDL(DATABASE, "create table values_table" +
columns));
+ assertTrue(setup.executeDDL(DATABASE, "create table empty_table" +
columns));
+ assertTrue(
+ setup.executeInsert(
+ DATABASE,
+ "insert into values_table values"
+ + "('nulls', null, null, null, null, null,
null, null, null, null),"
+ + "('zeros', false, 0, 0, 0, 0.0, 0.0, '',
'2020-02-29', 0),"
+ + "('values', true, -12, 42, 1234567890123,
1.5, 2.25, 'hello', '2024-01-02', 1704153600000);"));
+ assertTrue(
+ setup.executeDDL(
+ SECOND_DATABASE,
+ "create table orders(id string, amount int,
index(key=id));"));
+ assertTrue(setup.executeInsert(SECOND_DATABASE, "insert into orders
values('order', 7);"));
+ assertTrue(
+ setup.executeDDL(
+ DATABASE,
+ "create table bulk_table(id string, seq int,
index(key=id))"
+ + (Boolean.getBoolean("openmldb.cluster")
+ ? "
options(partitionnum=4,replicanum=1)"
+ : "")
+ + ";"));
+ StringBuilder bulk = new StringBuilder("insert into bulk_table
values");
+ for (int i = 0; i < 1101; i++) {
+ if (i > 0) {
+ bulk.append(',');
+ }
+ bulk.append("('key").append(i).append("',").append(i).append(')');
+ }
+ assertTrue(setup.executeInsert(DATABASE, bulk.append(';').toString()));
+ }
+
+ @AfterAll
+ static void cleanup() {
+ if (setup != null) {
+ setup.executeDDL(DATABASE, "drop table values_table;");
+ setup.executeDDL(DATABASE, "drop table empty_table;");
+ setup.executeDDL(DATABASE, "drop table bulk_table;");
+ setup.executeDDL(SECOND_DATABASE, "drop table orders;");
+ setup.dropDB(DATABASE);
+ setup.dropDB(SECOND_DATABASE);
+ setup.close();
+ }
+ }
+
+ @Test
+ void multiTablePreservesNullsTypesAndTableIdentity() throws Exception {
+ OpenMldbSource source =
+ source(
+ entry("empty", "select * from empty_table", FIELDS,
null)
+ + ","
+ + entry("values", "select * from
values_table", FIELDS, null)
+ + ","
+ + entry(
+ "orders",
+ "select amount as total, id as
order_id from orders",
+ "order_id=STRING, total=INT",
+ SECOND_DATABASE));
+ assertEquals(3, source.getProducedCatalogTables().size());
+ SingleSplitReaderContext context = context();
+ List<SeaTunnelRow> rows = new ArrayList<>();
+ try (OpenMldbSourceReader reader = (OpenMldbSourceReader)
source.createReader(context)) {
+ reader.open();
+ reader.pollNext(collector(rows));
+ reader.pollNext(collector(rows));
+ }
+ assertEquals(4, rows.size());
+ verify(context).signalNoMoreElement();
+ CatalogTable values = source.getProducedCatalogTables().get(1);
+ Map<String, SeaTunnelRow> byId = new HashMap<>();
+ for (SeaTunnelRow row : rows) {
+ if
(values.getTableId().toTablePath().toString().equals(row.getTableId())) {
+ byId.put((String) field(row, values.getSeaTunnelRowType(),
"id"), row);
+ }
+ }
+ assertEquals(3, byId.size());
+ SeaTunnelRow nulls = byId.get("nulls");
+ for (String name : new String[] {"b", "s", "i", "l", "f", "d", "text",
"day", "ts"}) {
+ assertNull(field(nulls, values.getSeaTunnelRowType(), name), name);
+ }
+ SeaTunnelRow zeros = byId.get("zeros");
+ assertEquals(false, field(zeros, values.getSeaTunnelRowType(), "b"));
+ assertEquals((short) 0, field(zeros, values.getSeaTunnelRowType(),
"s"));
+ assertEquals(0, field(zeros, values.getSeaTunnelRowType(), "i"));
+ assertEquals(0L, field(zeros, values.getSeaTunnelRowType(), "l"));
+ assertEquals(0F, field(zeros, values.getSeaTunnelRowType(), "f"));
+ assertEquals(0D, field(zeros, values.getSeaTunnelRowType(), "d"));
+ assertEquals("", field(zeros, values.getSeaTunnelRowType(), "text"));
+ assertEquals(LocalDate.of(2020, 2, 29), field(zeros,
values.getSeaTunnelRowType(), "day"));
+ assertEquals(
+ new Timestamp(0).toLocalDateTime(),
+ field(zeros, values.getSeaTunnelRowType(), "ts"));
+ SeaTunnelRow nonNull = byId.get("values");
+ Object[] expected = {
+ true,
+ (short) -12,
+ 42,
+ 1234567890123L,
+ 1.5F,
+ 2.25D,
+ "hello",
+ LocalDate.of(2024, 1, 2),
+ new Timestamp(1704153600000L).toLocalDateTime()
+ };
+ String[] names = {"b", "s", "i", "l", "f", "d", "text", "day", "ts"};
+ for (int i = 0; i < names.length; i++) {
+ assertEquals(
+ expected[i], field(nonNull, values.getSeaTunnelRowType(),
names[i]), names[i]);
+ }
+ assertNotSame(nulls.getFields(), zeros.getFields());
+ CatalogTable orders = source.getProducedCatalogTables().get(2);
+ SeaTunnelRow order =
+ rows.stream()
+ .filter(
+ row ->
+ row.getTableId()
+ .equals(
+ orders.getTableId()
+ .toTablePath()
+ .toString()))
+ .findFirst()
+ .get();
+ assertEquals("order", field(order, orders.getSeaTunnelRowType(),
"order_id"));
+ assertEquals(7, field(order, orders.getSeaTunnelRowType(), "total"));
+ }
+
+ @Test
+ void legacySourcePreservesNullsAndReaderOwnership() throws Exception {
+ OpenMldbSource legacy = new OpenMldbSource(parameters("select * from
values_table"));
+ assertTrue(
+ legacy.getProducedCatalogTables()
+ .get(0)
+ .getTableSchema()
+ .getColumns()
+ .get(1)
+ .isNullable());
+ SingleSplitReaderContext context = context();
+ try (OpenMldbSourceReader first = (OpenMldbSourceReader)
legacy.createReader(context);
+ OpenMldbSourceReader second = (OpenMldbSourceReader)
legacy.createReader(context)) {
+ first.open();
+ second.open();
+ first.close();
+ List<SeaTunnelRow> rows = new ArrayList<>();
+ second.pollNext(collector(rows));
+ assertEquals(3, rows.size());
+ SeaTunnelRow nulls =
+ rows.stream().filter(row ->
"nulls".equals(row.getField(0))).findFirst().get();
+ for (int i = 1; i < 10; i++) {
+ assertNull(nulls.getField(i));
+ }
+ }
+ }
+
+ @Test
+ void rejectsInvalidQueryAndSchemaWithoutSuccessfulCompletion() throws
Exception {
+ for (String entry :
+ new String[] {
+ entry("missing", "select * from missing_table",
"id=STRING", null),
+ entry("count", "select * from values_table", "id=STRING",
null),
+ entry("type", "select i from values_table", "i=STRING",
null),
+ entry("name", "select i from values_table", "wrong=INT",
null),
+ entry(
+ "duplicate",
+ "select i as value, i as value from values_table",
+ "value=INT, other=INT",
+ null)
+ }) {
+ OpenMldbSource source = source(entry);
+ SingleSplitReaderContext context = context();
+ List<SeaTunnelRow> rows = new ArrayList<>();
+ try (OpenMldbSourceReader reader =
+ (OpenMldbSourceReader) source.createReader(context)) {
+ reader.open();
+ OpenMldbConnectorException error =
+ assertThrows(
+ OpenMldbConnectorException.class,
+ () -> reader.pollNext(collector(rows)));
+ assertTrue(error.getCause() instanceof SQLException);
+ }
+ assertTrue(rows.isEmpty());
+ verify(context, never()).signalNoMoreElement();
+ }
+ }
+
+ @Test
+ void invalidQueryDoesNotPreventLaterQueriesOnSameClient() throws Exception
{
+ OpenMldbSource source =
+ source(entry("one", "select id from values_table",
"id=STRING", null));
+ SeaTunnelRowType type =
source.getProducedCatalogTables().get(0).getSeaTunnelRowType();
+ try (OpenMldbReadClient client =
+ new OpenMldbReadClient(parameters("select id from
values_table"))) {
+ assertThrows(
+ SQLException.class,
+ () -> client.execute(DATABASE, "select * from
missing_table", type, true));
+ try (OpenMldbReadClient.Query query =
+ client.execute(DATABASE, "select id from values_table",
type, true)) {
+ int rows = 0;
+ while (query.next()) {
+
assertFalse(query.readRow().getField(0).toString().isEmpty());
+ rows++;
+ }
+ assertEquals(3, rows);
+ }
+ }
+ }
+
+ private static Object field(SeaTunnelRow row, SeaTunnelRowType type,
String name) {
+ return row.getField(type.indexOf(name));
+ }
+
+ @Test
+ void readsAllRowsAcrossPartitions() throws Exception {
+ OpenMldbSource source =
+ source(entry("bulk", "select seq, id from bulk_table",
"id=STRING, seq=INT", null));
+ List<SeaTunnelRow> rows = new ArrayList<>();
+ try (OpenMldbSourceReader reader = (OpenMldbSourceReader)
source.createReader(context())) {
+ reader.open();
+ reader.pollNext(collector(rows));
+ }
+ assertEquals(1101, rows.size());
+ HashSet<Object> keys = new HashSet<>();
+ SeaTunnelRowType type =
source.getProducedCatalogTables().get(0).getSeaTunnelRowType();
+ for (SeaTunnelRow row : rows) {
+ Object id = field(row, type, "id");
+ assertEquals("key" + field(row, type, "seq"), id);
+ assertTrue(keys.add(id));
+ }
+ }
+
+ private static String connection() {
+ if (Boolean.getBoolean("openmldb.cluster")) {
+ return "cluster_mode=true\nzk_host=\""
+ + System.getProperty("openmldb.zk.host", "127.0.0.1:2181")
+ + "\"\nzk_path=\""
+ + System.getProperty("openmldb.zk.path", "/openmldb")
+ + "\"\ndatabase="
+ + DATABASE
+ + "\n";
+ }
+ return "cluster_mode=false\nhost=\""
+ + System.getProperty("openmldb.host", "127.0.0.1")
+ + "\"\nport="
+ + System.getProperty("openmldb.port", "6527")
+ + "\ndatabase="
+ + DATABASE
+ + "\n";
+ }
+
+ private static OpenMldbParameters parameters(String sql) {
+ return OpenMldbParameters.buildWithConfig(
+ ConfigFactory.parseString(connection() + "sql=\"" + sql +
"\""));
+ }
+
+ private static OpenMldbSource source(String entries) {
+ return (OpenMldbSource)
+ new OpenMldbSourceFactory()
+ .<SeaTunnelRow, SingleSplit,
SingleSplitEnumeratorState>createSource(
+ new TableSourceFactoryContext(
+ ReadonlyConfig.fromConfig(
+ ConfigFactory.parseString(
+ connection()
+ +
"tables_configs=["
+ + entries
+ + "]")),
+
TestOpenMldbSourceIT.class.getClassLoader()))
+ .createSource();
+ }
+
+ private static String entry(String table, String sql, String fields,
String database) {
+ return "{sql=\""
+ + sql
+ + "\", schema {table="
+ + table
+ + ", fields {"
+ + fields
+ + "}}"
+ + (database == null ? "" : ", database=" + database)
+ + "}";
+ }
+
+ private static SingleSplitReaderContext context() {
+ SingleSplitReaderContext context =
mock(SingleSplitReaderContext.class);
+ when(context.getBoundedness()).thenReturn(Boundedness.BOUNDED);
+ return context;
+ }
+
+ private static Collector<SeaTunnelRow> collector(List<SeaTunnelRow> rows) {
+ return new Collector<SeaTunnelRow>() {
+ @Override
+ public void collect(SeaTunnelRow row) {
+ rows.add(row);
+ }
+
+ @Override
+ public Object getCheckpointLock() {
+ return rows;
+ }
+ };
+ }
+}