跳到主要内容
版本:3.0.0

TDengine

TDengine 源端连接器

支持的引擎​

Spark
Flink
SeaTunnel Zeta

描述​

从 TDengine 超级表读取数据。

该 Source 以批处理方式按时间范围读取一个超级表的数据。可以读取该超级表下的所有子表,也可以只读取指定子表;同时支持只读取部分列。

每个 source 分片会读取一个 TDengine 子表。输出字段会自动把 subtable_name 放在第一列,用来保留原始子表名;如果下游继续写入 TDengine,TDengine sink 也会使用这一列作为目标子表名。

主要特性​

配置项​

名称类型必填默认值说明
urlString是-TDengine REST JDBC 连接地址,例如 jdbc:TAOS-RS://localhost:6041/。
usernameString是-连接 TDengine 使用的用户名。
passwordString是-连接 TDengine 使用的密码。
databaseString是-TDengine 数据库名称。
stableString是-TDengine 超级表名称。
sub_tablesList否-子表名称列表。不配置时读取该超级表下所有子表;配置后只读取指定子表。
lower_boundString是-查询时间范围的下界(包含),例如 2018-10-03 14:38:05.000。
upper_boundString是-查询时间范围的上界(不包含),例如 2018-10-03 14:38:16.801。
read_columnsList否-要读取的列名列表。不配置时读取所有列。读取超级表时请将 TAGS 列放在列表末尾;不要包含 subtable_name。
common-options否-Source 通用参数,请参考 Source Common Options。

url [String]​

TDengine REST JDBC 连接地址。

例如:

jdbc:TAOS-RS://localhost:6041/

username [String]​

连接 TDengine 使用的用户名。

password [String]​

连接 TDengine 使用的密码。

database [String]​

TDengine 数据库名称,数据库必须已经在服务端存在。

stable [String]​

TDengine 超级表名称。连接器会按该超级表下的每个子表生成一个 source 分片。

sub_tables [List]​

TDengine 子表名称列表。不配置时读取该超级表下的所有子表;配置后只读取指定子表。子表名需与服务端保持完全一致。

lower_bound [String]​

查询时间范围的下界,包含该时间点。连接器会为每个子表查询加上 timestamp_column >= lower_bound 条件。请使用 TDengine 可识别的时间字符串,例如 2018-10-03 14:38:05.000。

upper_bound [String]​

查询时间范围的上界,不包含该时间点。连接器会为每个子表查询加上 timestamp_column < upper_bound 条件。请使用 TDengine 可识别的时间字符串,例如 2018-10-03 14:38:16.801。

read_columns [List]​

要读取的列名列表。不配置时读取所有列。读取超级表时,请将 TAGS 字段放在列表末尾。不要在这里配置 subtable_name,连接器会自动把它作为第一列输出。

read_columns 的顺序会决定 subtable_name 之后的输出字段顺序。如果后续继续写入 TDengine Sink,请保持普通列在前、TAGS 字段在后,这样 Sink 才能正确区分普通列和 TAGS 值。

通用选项​

Source 插件通用参数,请参考 Source Common Options。

输出 Schema​

输出表的第一列固定为预留字段 subtable_name,用来标识该行来自哪个 TDengine 子表。后续字段为 read_columns 中声明的列(未设置时为所有列),顺序与 read_columns 保持一致;TAGS 列请按声明的顺序排在普通列之后。

示例​

按时间范围读取所有子表​

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

source {
TDengine {
url = "jdbc:TAOS-RS://localhost:6041/"
username = "root"
password = "taosdata"
database = "power"
stable = "meters"
lower_bound = "2018-10-03 14:38:05.000"
upper_bound = "2018-10-03 14:38:16.801"
plugin_output = "tdengine_result"
}
}

读取指定子表和指定列​

source {
TDengine {
url = "jdbc:TAOS-RS://localhost:6041/"
username = "root"
password = "taosdata"
database = "power"
stable = "meters"
lower_bound = "2018-10-03 14:38:05.000"
upper_bound = "2018-10-03 14:38:16.801"
sub_tables = ["d1001", "d1002"]
read_columns = ["ts", "current", "voltage", "phase", "off", "nc", "location", "groupid"]
}
}

该配置只会读取超级表 meters 下的子表 d1001 和 d1002。TAGS 字段(location、groupid)放在 read_columns 末尾,方便下游 TDengine Sink 正确切分普通列与 TAGS 值。

从 TDengine 读取并写入 TDengine​

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

source {
TDengine {
url = "jdbc:TAOS-RS://tdengine-src:6041/"
username = "root"
password = "taosdata"
database = "power"
stable = "meters"
lower_bound = "2018-10-03 14:38:05.000"
upper_bound = "2018-10-03 14:38:16.801"
plugin_output = "tdengine_result"
}
}

sink {
TDengine {
url = "jdbc:TAOS-RS://tdengine-sink:6041/"
username = "root"
password = "taosdata"
database = "power2"
stable = "meters2"
timezone = "UTC"
}
}

写入时,Sink 会使用源端产生的 subtable_name 作为目标子表名。TAGS 字段的数量会从目标超级表的元数据中读取,因此目标超级表必须先在服务端创建好。

变更日志​

Change Log
ChangeCommitVersion
[Improve][Connector-V2] Improve source split round-robin assignment for InfluxDB IoTDB Kudu and TDengine (#11573)https://github.com/apache/seatunnel/commit/1fe15cbb83.0.0