Skip to main content
Version: Next

OpenMldb

OpenMldb source connector

Support Those Engines​

Spark
Flink
SeaTunnel Zeta

Description​

Used to read data from OpenMLDB. The connector executes the configured SQL statement against 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​

Data Type Mapping​

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 TypeSeaTunnel Data Type
boolboolean
smallintsmallint
intint
bigintbigint
float / doublefloat / double
string / varcharstring
datedate
timestamptimestamp

Source Options​

nametyperequireddefault valuedescription
cluster_modebooleanyes-Whether to connect to OpenMLDB in cluster mode. Set to false for standalone mode.
sqlstringconditional-Single-query SQL. Configure either this option or tables_configs, not both.
tables_configslistconditional-Queries and explicit result schemas for a multi-table read.
databasestringyes-The OpenMLDB database name to connect to.
hoststringno-Required when cluster_mode is false. Host of the standalone OpenMLDB server.
portintno-Required when cluster_mode is false. Port of the standalone OpenMLDB server.
zk_hoststringno-Required when cluster_mode is true. ZooKeeper host list of the OpenMLDB cluster.
zk_pathstringno-Required when cluster_mode is true. ZooKeeper path of the OpenMLDB cluster.
session_timeoutintno10000OpenMLDB session timeout in milliseconds.
request_timeoutintno60000OpenMLDB request timeout in milliseconds.
common-optionsno-Source plugin common parameters, please refer to Source Common Options for details.

cluster_mode [boolean]​

Whether to connect to OpenMLDB in cluster mode. When it is false, configure host and port. When it is true, configure zk_host and zk_path.

sql [string]​

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.

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]​

The OpenMLDB database name to connect to. The configured database must exist on the target OpenMLDB instance.

host [string]​

OpenMLDB host. Only used when cluster_mode is false (standalone mode).

port [int]​

OpenMLDB port. Only used when cluster_mode is false (standalone mode).

zk_host [string]​

ZooKeeper host list for the OpenMLDB cluster, for example zk-1:2181,zk-2:2181,zk-3:2181. Only used when cluster_mode is true.

zk_path [string]​

ZooKeeper path of the OpenMLDB cluster, for example /openmldb. Only used when cluster_mode is true.

session_timeout [int]​

OpenMLDB session timeout in milliseconds. Defaults to 10000 (10 seconds).

request_timeout [int]​

OpenMLDB request timeout in milliseconds. Defaults to 60000 (60 seconds).

common options​

Source plugin common parameters, please refer to Source Common Options for details.

Task Example​

Standalone mode​

source {
OpenMldb {
host = "172.17.0.2"
port = 6527
sql = "select * from demo_table1"
database = "demo_db"
cluster_mode = false
}
}

Cluster mode​

source {
OpenMldb {
zk_host = "zk-1:2181,zk-2:2181,zk-3:2181"
zk_path = "/openmldb"
sql = "select * from demo_table1"
database = "demo_db"
cluster_mode = true
}
}

With downstream sink​

A typical end-to-end job that reads from OpenMLDB and prints the rows through the Console sink.

env {
parallelism = 1
job.mode = "BATCH"
}

source {
OpenMldb {
host = "172.17.0.2"
port = 6527
sql = "select id, name from demo_table1"
database = "demo_db"
cluster_mode = false
}
}

sink {
Console {
}
}

Multi-table read​

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​

Change Log
ChangeCommitVersion
[improve] openmldb options (#9166)https://github.com/apache/seatunnel/commit/d324fc59a42.3.11
[Improve][Common] Introduce new error define rule (#5793)https://github.com/apache/seatunnel/commit/9d1b2582b22.3.4
[Improve][build] Give the maven module a human readable name (#4114)https://github.com/apache/seatunnel/commit/d7cd6010512.3.1
[Improve][Project] Code format with spotless plugin. (#4101)https://github.com/apache/seatunnel/commit/a2ab1665612.3.1
[Feature][Connector] add get source method to all source connector (#3846)https://github.com/apache/seatunnel/commit/417178fb842.3.1
[Feature][API & Connector & Doc] add parallelism and column projection interface (#3829)https://github.com/apache/seatunnel/commit/b9164b8ba12.3.1
[Hotfix][OptionRule] Fix option rule about all connectors (#3592)https://github.com/apache/seatunnel/commit/226dc6a1192.3.0
[Feature][Conenctor-V2] Add openmldb source connector (#3313)https://github.com/apache/seatunnel/commit/e68ecf7bef2.3.0