零代码数据接入指南:基于 taosExplorer 与 taosX 的数据源集成、ETL 与任务管理
2026/9/13 9:45:05 网站建设 项目流程

零代码数据接入指南:基于 taosExplorer 与 taosX 的数据源集成、ETL 与任务管理

【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine

本文档基于 TDengine 官方指南《Data Connectors》(docs/en/08-data-ingest-and-delivery/01-no-code-ingestion/index.md)编写,系统讲解如何通过 taosExplorer 与 taosX 以零代码方式将第三方数据源持续写入 TDengine。

引言

TDengine 内置了可视化数据管理工具taosExplorer(即 TDengine TSDB Explorer)。借助 taosExplorer,用户只需在浏览器中通过简单配置即可向 TDengine 提交任务,实现从多种数据源到 TDengine 的零代码数据导入。导入过程中,TDengine 自动完成数据的解析、过滤和转换,保证导入数据质量。通过这种零代码数据源集成方式,TDengine 已成为时序大数据聚合平台:用户无需额外部署 ETL 工具,从而极大简化整体架构设计、提升数据处理效率。

从架构上看,零代码集成平台由两个核心组件构成:

  • taosExplorer:提供可视化 Web 界面,用于创建数据接入任务、配置解析规则、预览结果、监控任务状态。
  • taosX:负责连接数据源、执行解析/过滤/转换逻辑并写入 TDengine 的服务端组件。

当 taosX 无法直连数据源时(例如数据源位于隔离的 OT 网络、工业控制网段或受限内网),可在数据源所在网络部署taosX-Agent作为代理,由 Agent 采集数据并转发给 taosX,详见安装 taosX-Agent与taosX-Agent 参考。

图:零代码集成平台的系统架构。

支持的数据源

TDengine 当前支持的数据源如下:

数据源支持版本说明
Aveva PI SystemPI AF Server 2.10.9.593 及以上工业数据管理与分析平台(前身为 OSIsoft PI System),可实时采集、集成、分析与可视化工业数据,帮助企业实现智能决策与精细化管理
Aveva HistorianAVEVA Historian 2020 RS SP1工业大数据分析软件(前身为 Wonderware Historian),面向工业环境存储、管理和分析来自各类工业设备与传感器的实时及历史数据
OPC DAMatrikon OPC 1.7.2.7433Open Platform Communications 缩写,一种开放、标准化的通信协议,用于不同制造商自动化设备之间的数据交换。最初由微软开发,旨在解决工业控制领域的互操作性问题;OPC 协议于 1996 年首次发布,当时称为 OPC DA(Data Access),主要用于实时数据采集与控制
OPC UAKeepWare KEPServerEx 6.52006 年 OPC 基金会发布 OPC UA(Unified Architecture)标准,面向服务、面向对象,灵活性和可扩展性更高,已成为 OPC 协议的主流版本
MQTTemqx: 3.0.0 ~ 5.7.1;hivemq: 4.0.0 ~ 4.31.0;mosquitto: 1.4.4 ~ 2.0.18Message Queuing Telemetry Transport 缩写,基于发布/订阅模式的轻量级通信协议,低开销、低带宽占用,广泛适用于物联网、小型设备、移动应用等领域
Kafka2.11 ~ 3.8.0Apache 软件基金会开发的开源流处理平台,主要用于处理实时数据,提供统一、高吞吐、低延迟的消息系统。具有高速、可扩展、持久化与分布式设计,可每秒处理数十万次读写、支持数千客户端,同时保持数据可靠性与可用性
InfluxDB1.7、1.8、2.0-2.7广受欢迎的开源时序数据库,针对海量时序数据的处理进行了优化
OpenTSDB2.4.1基于 HBase 的分布式、可扩展时序数据库,主要用于存储、索引和访问从大规模集群(包括网络设备、操作系统、应用等)采集的指标数据,使数据更易于访问和图形化展示
MySQL5.6、5.7、8.0+最流行的关系型数据库管理系统之一,以体积小、速度快、总体拥有成本低著称,特别是其开源特性,使其成为中大型网站数据库开发的选择
Oracle11G/12c/19c世界流行的关系型数据库管理系统之一,可移植性好、使用方便、功能强大,适用于各种大中小型计算机环境,是高效、可靠、高吞吐的数据库解决方案
PostgreSQLv15.0+非常强大的开源客户端/服务器关系型数据库管理系统,具备大型商业 RDBMS 的众多特性,包括事务、子查询、触发器、视图、外键引用完整性和复杂锁定能力
SQL Server2012/2022微软开发的关系型数据库管理系统,易用性好、可扩展性强、与相关软件集成度高
MongoDB3.6+介于关系型与非关系型数据库之间的产品,广泛用于内容管理系统、移动应用和物联网等众多领域
CSV-Comma Separated Values 缩写,以逗号分隔的纯文本文件格式,常用于电子表格或数据库软件
TDengine Query2.4+、3.0+从旧版本 TDengine 查询并写入新集群
TDengine Data Subscription3.0+使用 TMQ 订阅 TDengine 中指定的数据库或超级表

数据抽取、过滤与转换(内置 ETL)

由于数据源众多,各数据源可能具有不同的物理单位、命名约定和时区。为应对这一问题,TDengine 内置了 ETL 能力,可以从数据源的数据包中解析并提取所需数据,并执行过滤与转换,以保证写入数据的质量并提供统一的命名空间。具体功能如下:

  1. 解析:使用 JSON Path 或正则表达式从原始消息中解析字段。
  2. 从列中抽取或拆分:使用 split 或正则表达式从一个原始字段中提取多个字段。
  3. 过滤:仅当表达式值为真时,消息才写入 TDengine。
  4. 转换:在解析字段与 TDengine 超级表字段之间建立转换与映射关系。

解析(Parsing)

仅非结构化数据源需要此步骤。目前 MQTT 和 Kafka 数据源使用本步骤提供的规则解析非结构化数据,以初步获得结构化数据(即可用字段描述的"行列"数据)。在 Explorer 中,需要提供示例数据和解析规则,并预览解析后的结构化数据表格。

示例数据(Sample Data)

图:示例数据输入区。

文本区域中存放示例数据,可通过三种方式获取:

  1. 直接在文本区域中输入示例数据;
  2. 点击右侧"Retrieve from Server"(从服务器获取)按钮,从已配置的服务器获取示例数据并追加到示例数据文本区域;
  3. 上传文件,将文件内容追加到示例数据文本区域。

每条示例数据以回车符结尾。

解析规则

解析是将非结构化字符串解析为结构化数据的过程。消息体的解析规则目前支持JSONRegex(正则表达式)UDT(自定义解析脚本)三种。

1. JSON 解析

JSON 解析支持 JSONObject 或 JSONArray。以下 JSON 示例数据可自动解析出字段:groupidvoltagecurrenttsinuselocation

{"groupid": 170001, "voltage": "221V", "current": 12.3, "ts": "2023-12-18T22:12:00", "inuse": true, "location": "beijing.chaoyang.datun"} {"groupid": 170001, "voltage": "220V", "current": 12.2, "ts": "2023-12-18T22:12:02", "inuse": true, "location": "beijing.chaoyang.datun"} {"groupid": 170001, "voltage": "216V", "current": 12.5, "ts": "2023-12-18T22:12:04", "inuse": false, "location": "beijing.chaoyang.datun"}

[{"groupid": 170001, "voltage": "221V", "current": 12.3, "ts": "2023-12-18T22:12:00", "inuse": true, "location": "beijing.chaoyang.datun"}, {"groupid": 170001, "voltage": "220V", "current": 12.2, "ts": "2023-12-18T22:12:02", "inuse": true, "location": "beijing.chaoyang.datun"}, {"groupid": 170001, "voltage": "216V", "current": 12.5, "ts": "2023-12-18T22:12:04", "inuse": false, "location": "beijing.chaoyang.datun"}]

后续示例仅以 JSONObject 说明。以下嵌套 JSON 数据可自动解析出字段groupiddata_voltagedata_currenttsinuselocation_0_provincelocation_0_citylocation_0_datun;同时可以选择要解析的字段,并为解析出的字段设置别名。

{"groupid": 170001, "data": { "voltage": "221V", "current": 12.3 }, "ts": "2023-12-18T22:12:00", "inuse": true, "location": [{"province": "beijing", "city":"chaoyang", "street": "datun"}]}

图:嵌套 JSON 的解析与字段别名设置。

2. Regex 正则表达式

可以使用正则表达式中的命名捕获组(named capture groups)从任意字符串(文本)字段中提取多个字段。例如下图从 nginx 日志中提取访问 IP、时间戳、访问 URL 等字段:

(?<ip>\b(?:[0-9]{1,3}\.){3}[0-9]{1,3}\b)\s-\s-\s\[(?<ts>\d{2}/\w{3}/\d{4}:\d{2}:\d{2}:\d{2}\s\+\d{4})\]\s"(?<method>[A-Z]+)\s(?<url>[^\s"]+).*(?<status>\d{3})\s(?<length>\d+)

图:使用命名捕获组解析 nginx 日志。

3. UDT 自定义解析脚本

支持编写自定义 rhai 语法脚本解析输入数据(语法参考 rhai 官方文档),脚本目前仅支持 json 格式的原始数据。

  • 输入:脚本中可使用参数data,它是原始数据经 json 解析后的 Object Map;
  • 输出:输出数据必须是一个数组。

例如,某设备上报三相电压值,需要分别写入三张子表,待解析的数据如下:

{ "ts": "2024-06-27 18:00:00", "voltage": "220.1,220.3,221.1", "dev_id": "8208891" }

可使用如下脚本提取三路电压数据:

let v3 = data["voltage"].split(","); [ #{"ts": data["ts"], "val": v3[0], "dev_id": data["dev_id"]}, #{"ts": data["ts"], "val": v3[1], "dev_id": data["dev_id"]}, #{"ts": data["ts"], "val": v3[2], "dev_id": data["dev_id"]} ]

最终解析结果如下图所示:

图:UDT 脚本解析三相电压数据的结果。

抽取或拆分(Extraction or Splitting)

解析后的数据可能仍不满足目标表的数据要求。例如智能电表采集的原始数据如下(json 格式):

{"groupid": 170001, "voltage": "221V", "current": 12.3, "ts": "2023-12-18T22:12:00", "inuse": true, "location": "beijing.chaoyang.datun"} {"groupid": 170001, "voltage": "220V", "current": 12.2, "ts": "2023-12-18T22:12:02", "inuse": true, "location": "beijing.chaoyang.datun"} {"groupid": 170001, "voltage": "216V", "current": 12.5, "ts": "2023-12-18T22:12:04", "inuse": false, "location": "beijing.chaoyang.datun"}

使用 json 规则解析后,voltage 被解析为带单位的字符串,而用户希望以 int 类型记录电压和电流值用于统计分析,因此需要进一步拆分 voltage;此外,还希望将日期拆分为日期和时间分别存储。

可对源字段ts使用 split 规则拆分为日期和时间,并对字段voltage使用 regex 提取电压值与单位。split 规则需要设置分隔符(delimiter)拆分次数(number of splits),拆分字段的命名规则为{原字段名}_{序号}。Regex 规则与解析过程相同,使用命名捕获组命名提取出的字段。

过滤(Filtering)

过滤功能可设置过滤条件,只有满足条件的数据行才会写入目标表。过滤条件表达式的结果必须为布尔类型。编写过滤条件前,需要确定解析字段的类型,并根据字段类型使用判断函数和比较运算符(>>=<=<==!=)进行判断。

字段类型与转换

只有明确解析出的每个字段的类型,才能使用正确的语法进行数据过滤。

使用 json 规则解析出的字段会根据属性值自动设置类型:

  1. bool 类型:"inuse": true
  2. int 类型:"voltage": 220
  3. float 类型:"current" : 12.2
  4. string 类型:"location": "MX001"

使用 regex 规则解析出的数据均为 string 类型;使用 split 和 regex 抽取或拆分出的数据为 string 类型。

如果提取出的数据类型不是预期类型,可以进行数据类型转换。常见的数据类型转换是将字符串转换为数值类型。支持的转换函数如下:

函数从类型到类型示例
parse_intstringintparse_int("56") // 结果为整数 56
parse_floatstringfloatparse_float("12.3") // 结果为浮点数 12.3
条件表达式

不同数据类型有其各自的条件表达式写法。

1. BOOL 类型

可以使用变量本身或!运算符。例如对于字段 "inuse": true,可写如下表达式:

  1. inuse
  2. !inuse

2. 数值类型(int/float)

数值类型支持比较运算符==!=>>=<<=

3. 字符串类型

使用比较运算符比较字符串。字符串函数如下:

函数说明示例
is_empty字符串为空时返回 trues.is_empty()
contains检查字符串中是否包含某字符或子串s.contains("substring")
starts_with字符串以某字符串开头时返回 trues.starts_with("prefix")
ends_with字符串以某字符串结尾时返回 trues.ends_with("suffix")
len返回字符串的字符数(非字节数),必须与比较运算符配合使用s.len == 5 检查字符串长度是否为 5;len 作为属性返回 int,与前四个直接返回 bool 的函数不同

4. 复合表达式

多个条件表达式可使用逻辑运算符(&&、||、!)组合。例如以下表达式表示获取安装在北京且电压值大于 200 的智能电表数据:

location.starts_with("beijing") && voltage > 200

映射(Mapping)

映射是将解析、抽取或拆分出的源字段映射到目标表字段。可以直接映射,也可以经过一定规则计算后再映射到目标表。

选择目标超级表

选择目标超级表后,将加载该超级表的所有标签(tag)和列(column)。源字段会根据名称按映射规则自动映射到目标超级表的标签和列。

映射规则

支持的映射规则如下表:

规则说明
mapping直接映射,需要选择映射源字段
value常量,可输入字符串常量或数值常量,输入的常量值直接存储
generator生成器,目前仅支持时间戳生成器 now,存储时写入当前时间
join字符串连接器,可指定连接字符拼接选中的多个源字段
format字符串格式化工具,填写格式化字符串,例如有三个源字段 year、month、day 分别表示年、月、日,希望以 yyyy-MM-dd 日期格式存储,可提供格式化字符串${year}-${month}-${day}。其中${}作为占位符,占位符可以是源字段或字符串类型字段函数处理
sum选择多个数值字段进行加法计算
expr数值运算表达式,可对数值字段进行更复杂的函数处理和数学运算

1.format中支持的字符串处理函数

函数说明示例
pad(len, pad_chars)用字符或字符串将字符串填充到至少指定长度"1.2".pad(5, '0') // 结果为 "1.200"
trim去除字符串开头和结尾的空白" abc ee ".trim() // 结果为 "abc ee"
sub_string(start_pos, len)提取子字符串,两个参数:1. 起始位置,< 0 时从末尾计数;2.(可选)要提取的字符数,≤ 0 时不提取,省略时提取到末尾"012345678".sub_string(5) // "5678";"012345678".sub_string(5, 2) // "56";"012345678".sub_string(-2) // "78"
replace(substring, replacement)用另一字符串替换子串"012345678".replace("012", "abc") // "abc345678"

2.expr中的数学表达式

基本数学运算支持加法+、减法-、乘法*、除法/。例如数据源采集的温度值为摄氏度,而目标数据库存储值为华氏度,则需要对采集的温度数据进行换算。若源字段为temperature,则使用表达式temperature * 1.8 + 32

数学表达式还支持使用数学函数,如下表:

函数说明示例
sin, cos, tan, sinh, cosh三角函数a.sin()
asin, acos, atan, asinh, acosh反三角函数a.asin()
sqrt平方根a.sqrt() // 4.sqrt() == 2
exp指数a.exp()
ln, log对数a.ln() // e.ln() == 1;a.log() // 10.log() == 1
floor, ceiling, round, int, fraction取整a.floor() // (4.2).floor() == 4;a.ceiling() // (4.2).ceiling() == 5;a.round() // (4.2).round() == 4;a.int() // (4.2).int() == 4;a.fraction() // (4.2).fraction() == 0.2
子表名映射

子表名是字符串,可以使用映射规则中的字符串格式化format表达式定义。

创建任务(以 MQTT 为例)

下面以 MQTT 数据源为例,说明如何创建 MQTT 类型的任务,从 MQTT Broker 消费数据并写入 TDengine。

  1. 登录 taosExplorer 后,点击左侧导航栏中的 "Data Writing"(数据写入)进入任务列表页。
  2. 在任务列表页点击 "+ Add Data Source"(添加数据源)进入任务创建页。
  3. 输入任务名称后,选择类型为 MQTT,然后可以新建代理或选择已创建的代理。
  4. 输入 MQTT broker 的 IP 地址和端口号,例如:192.168.1.100:1883。
  5. 配置认证与 SSL 加密:
    • 如果 MQTT broker 已启用用户认证,在认证区域输入 MQTT broker 的用户名和密码;
    • 如果 MQTT broker 已启用 SSL 加密,可打开页面上的 SSL 证书开关,上传 CA 证书,以及客户端的证书和私钥文件。
  6. 在 "Collection Configuration"(采集配置)区域,可选择 MQTT 协议版本,目前支持 3.1、3.1.1、5.0;配置 Client ID 时注意,如果对同一 MQTT broker 创建多个任务,Client ID 应各不相同以避免冲突,否则可能导致任务无法正常运行;配置 topic 和 QoS 时使用<topic name>::<QoS>格式,QoS 取值范围为 0、1、2,分别表示至多一次、至少一次、恰好一次;配置完上述信息后可点击 "Check Connectivity"(检查连通性)按钮检查配置,若连通性检查失败,请根据页面返回的具体错误提示进行修改。
  7. 在从 MQTT broker 同步数据的过程中,taosX 还支持对消息体中的字段进行提取、过滤和映射操作。在 "Payload Transformation"(消息体转换)下的文本框中可直接输入消息体样本,或通过上传文件导入,未来还将支持直接从已配置的服务器获取样本消息。
  8. 对于消息体字段的提取,目前支持 JSON 和正则表达式两种方式。对于简单的 key/value 格式 JSON 数据,可直接点击提取按钮展示解析出的字段名;对于复杂 JSON 数据,可使用 JSON Path 提取感兴趣的字段;使用正则表达式提取字段时,需确保正则表达式的正确性。
  9. 消息体字段解析完成后,可根据解析出的字段名设置过滤规则,只有符合过滤规则的数据才会写入 TDengine,否则消息将被忽略。例如可配置过滤规则为 voltage > 200,表示只有电压大于 200V 的数据才会同步到 TDengine。
  10. 最后,配置消息体字段与超级表字段之间的映射规则后即可提交任务。除了基本映射外,这里还可以对消息体字段的值进行转换,例如使用表达式(expr)从原始消息体的电压和电流计算出功率后再写入 TDengine。
  11. 提交任务后将自动返回任务列表页,如果提交成功,任务状态将切换为 "Running";如果提交失败,可查看任务的运行日志(activity log)寻找错误原因。
  12. 对于正在运行的任务,点击指标(metrics)的查看按钮可查看任务的详细运行指标,弹窗分为 2 个标签页,分别显示任务多次运行的累计指标和本次运行的指标,这些指标每 2 秒自动刷新一次。

MQTT 连接的更多细节

结合 MQTT 数据源完整配置指南,连接与采集配置还包含以下要点:

  • TLS 校验有三种模式:Disabled(不校验 TLS 证书,连接器先尝试 TCP,失败后尝试不带证书校验的 TLS)、单向认证(使用 TLS 并校验服务器证书,上传 CA 证书)、双向认证(使用双向 TLS,上传 CA 证书、客户端证书和客户端私钥)。
  • Client ID:输入标识符后,将生成带taosx前缀的 client id(例如输入foo生成taosxfoo)。若开启末尾开关,则会在taosx与输入标识符之间拼接当前任务的任务 id(生成形如taosx100foo的 client id)。连接同一 MQTT 地址的所有 client id 必须唯一。
  • Keep Alive:若 broker 在 keep alive 时间内未收到客户端消息,将认为客户端已断开并关闭连接。
  • Clean Session:选择是否清除会话,默认值为 true。
  • Topics QoS Config:格式为{topic_name}::{qos}(如my_topic::0)。MQTT 5.0 支持共享订阅,多个客户端可订阅同一 topic 实现负载均衡,格式为$share/{group_name}/{topic_name}::{qos},其中$share是固定前缀,group_name是客户端组名,类似 Kafka 的 consumer group。
  • Topic Analysis(主题解析):解析 MQTT Topic 的每一级为对应变量名,_表示解析时忽略当前级。例如 MQTT Topica/+/c对应解析规则v1/v2/_,表示将第一级a赋给变量v1,第二级值(通配符+代表任意值)赋给变量v2,第三级c被忽略。在后续的消息体解析中,主题解析得到的变量也可参与各种转换和计算。
  • Compression(压缩):配置消息体压缩算法,taosX 接收消息后按对应算法解压。选项包括 none(不压缩)、gzip、snappy、lz4、zstd,默认 none。
  • Char Encoding(字符编码):配置消息体编码格式,选项包括 UTF_8、GBK、GB18030、BIG5,默认 UTF_8。

高级选项(Advanced Options)

Advanced Options区域默认折叠,点击>展开。MQTT 和 SparkplugB 数据源通常提供以下选项(字段名可能因连接器而异):

  • Message Queue Size(消息队列大小):指定接收缓冲区大小。如果队列已满且Cache Realtime Data未启用,新到达的数据将被丢弃。设置为0可禁用缓冲。
  • Maximum In-Process Batches(最大在途批次):指定可并发处理的批次数。达到该限制后,连接器停止从接收队列取消息,导致消息积压。最小值为1
  • Batch Size(批次大小):指定一次送入处理管道的消息数。与Batch Delay配合:即使延迟未到,批次满也会立即发送。最小值为1
  • Batch Delay(批次延迟):指定每批的毫秒级超时,从该批第一条消息起算。超时后即使未达到Batch Size也会发送。最小值为1
  • Write Concurrency(写入并发):指定可并发写入 TDengine 的任务数。
  • Cache Realtime Data(缓存实时数据):启用后,消费的数据先写入本地文件,由后台任务转发给下游。当下游处理跟不上时提供流量整形(traffic shaping)。积压数据消费完后缓存文件被删除。默认关闭。详见 Store and Forward。
  • Cache Storage Directory(缓存存储目录):覆盖缓存文件目录。仅在启用Cache Realtime Data时生效,否则默认使用 taosX 启动时配置的数据目录。
  • 启用Save Raw Data(保存原始数据)后,还可配置Maximum Retention Days(最大保留天数)Raw Data Storage Directory(原始数据存储目录)

异常处理策略(Exception Handling)

Exception Handling Strategy区域默认折叠,点击>展开。常用策略包括:

  • Archive(归档):将无效数据写入归档文件(默认位于${data_dir}/tasks/<id>/<datetime>下),不写入目标数据库。
  • Discard(丢弃):忽略无效数据。
  • Error(报错):报告错误。
  • Cache(缓存):针对目标连接失败或资源不足,将数据写入缓存文件,待目标恢复后再摄入。

可为以下场景配置策略:

场景可选策略
目标连接超时Archive、Discard、Error、Cache
目标数据库不存在Archive、Discard、Error
表不存在Archive、Discard、Error、自动建表并重试
主时间戳超出范围(now - keep1now + 100yArchive、Discard、Error
主时间戳为 nullArchive、Discard、Error、使用当前时间
复合主键为 nullArchive、Discard、Error
表名超过 192 字符Archive、Discard、Error、截断、截断并归档
表名包含非法字符(如.Archive、Discard、Error、用配置的字符串替换非法字符
表名模板变量为 nullDiscard、变量留空、用配置的字符串替换
列不存在Archive、Discard、Error、自动补列并重试
列名超过 64 字符Archive、Discard、Error
列值超过定义长度Archive、Discard、Error、截断、截断并归档;Automatic Column Expansion可改为修改表结构后重试
其他数据错误Archive、Discard、Error

其他附加设置:

  • Connection Timeout(连接超时):目标连接超时秒数,范围1~600
  • Temporary Storage Location(临时存储位置):相对${data_dir}/tasks/<id>/的路径。
  • Archive Retention Days(归档保留天数):非负整数,0表示不限。
  • Archive Available Space(归档可用空间)0~655350表示不限。
  • Archive Location(归档位置):相对${data_dir}/tasks/<id>/的路径。
  • Archive Write Failure Strategy(归档写入失败策略):删除旧文件、丢弃数据、或报错并停止任务。

任务管理

在任务列表页,可以对任务执行启动、停止、查看、删除、复制等操作,还可以查看每个任务的运行状态,包括写入记录数、流量等。

健康状态(Health Status)

从 v3.3.5.0 开始,任务列表为每个运行中的任务显示健康状态。Advanced Options区域提供以下健康监控设置:

  1. Health Check Duration(健康检查时长):计算任务状态的最近时间段。
  2. Busy State Threshold(繁忙状态阈值):排队项与写入队列容量的比值,默认100%
  3. Max Write Queue Length(最大写入队列长度):写入队列的最大容量。
  4. Write Error Threshold(写入错误阈值):健康检查期间允许的写入错误数,超过即报错。

任务列表可显示以下状态:

  • Ready(就绪):源与目标健康检查均通过,可正常读写。
  • Idle(空闲):监控周期内没有数据进入处理管道。
  • Active(活跃):任务正常运行并处理数据。
  • Pending(挂起):源健康,但写入端正在等待且没有消息可写。
  • Busy(繁忙):写入队列超过配置阈值。这可能表明存在性能瓶颈,需要调整参数或资源,但本身并不代表错误。
  • Bounce(抖动):源与目标健康,但写入错误超过阈值。这可能表明存在大量无效数据或数据丢失。
  • SourceError(源错误):源无法读取,工作负载会尝试重连。
  • SinkError(目标错误):目标无法写入,工作负载会尝试重连,恢复后回到Ready
  • Fatal(致命):发生严重或不可恢复的错误。

健康状态为空表示任务尚无数据进入。

从检查点恢复任务(Checkpoint 断点续传)

大多数 taosX 数据源可以从最后持久化的检查点恢复:

  • TDengine Query、MySQL、PostgreSQL、Oracle、Microsoft SQL Server、MongoDB:持久化上次查询的时间戳并从该时间点恢复。
  • TDengine Data Subscription:依赖 TDengine 订阅进度。
  • PI、InfluxDB、OpenTSDB:持久化消息时间戳并从该时间点恢复。
  • OPC UA、OPC DA、MQTT:启用Cache Realtime Data时将消息持久化到磁盘。
  • Kafka:依赖 Kafka 的 consumer-offset 管理。
  • CSV:持久化文件名与已消费行数,从下一行继续。
  • AVEVA Historian:查询历史数据时持久化最后一条消息的时间戳。
  • SparkplugB:目前不支持消息持久化与恢复。

进阶实践:完整的零代码接入流程

如果你希望从一个端到端的视角快速跑通零代码接入,可以参考快速入门章节 Zero-Code Data Ingestion 中的完整流程(以 MQTT 为例):

  1. 准备环境:确认 TDengine 服务运行、taosExplorer 可在浏览器访问、taosX 相关服务运行且 Explorer 左侧菜单包含Data In入口,并能访问外部 MQTT broker(该示例使用 EMQX 提供的公共 brokerbroker.emqx.io:1883)。
  2. 创建任务:在 Explorer 中进入Data In,点击+ Add Source新建任务,设置名称(如quick_mqtt_meter)、类型为MQTT、目标数据库(如test_mqtt)。
  3. 配置连接:填写 MQTT Host 与端口、协议版本(3.1 / 3.1.1 / 5.0)、Client ID(可留空使用自动生成)、Keep Alive、Clean Session、认证信息与 TLS 选项,并将 Topics QoS Config 配置为tdengine/quickstart/meter::0(格式为<topic>::<QoS>,多个 topic 用逗号分隔)。随后点击Check Connectivity验证连通性。
  4. 配置消息体转换:在 Payload Transformation 中输入 JSON 样例(如智能电表读数),点击识别/预览按钮确认 JSON 字段被正确提取。
  5. 配置映射:选择或创建目标超级表meters(配置TIMESTAMP类型的ts列、DOUBLE类型的current/phase列、INT类型的voltage列,以及groupidlocation标签),将 SubTableName 设置为t_{id}以按消息中的id生成子表名,并完成其余字段的映射后提交任务。
  6. 发送测试数据并验证:任务状态变为Running后,使用mosquitto_pub向同一 topic 发布消息:
mosquitto_pub -h broker.emqx.io -p 1883 -t 'tdengine/quickstart/meter' -q 0 -m '{ "ts": "2026-07-27T14:31:00+08:00", "id": 1, "current": 10.58, "phase": 1.41, "voltage": 221, "groupid": 7, "location": "beijing" }'

参数说明:-h/-p指定 broker 地址与端口;-t指定发布 topic(须与 Topics QoS Config 中的 topic 一致,不带::QoS后缀);-q指定 QoS;-m指定消息体,字段须与样例和映射规则一致。重复测试时请修改ts为当前时间以避免覆盖同一时间戳。

然后在 Explorer 的数据浏览器或 shell 中查询验证:

SELECT tbname, ts, current, voltage, phase, groupid, location FROM test_mqtt.meters ORDER BY ts DESC LIMIT 5;

若返回类似t_1 | 2026-07-27 14:31:00.000 | 10.5800 | 221 | 1.410 | 7 | beijing的行,说明 MQTT 消息已通过 taosX 写入 TDengine。

  1. 故障排查:若连通性检查失败,检查 taosX 所在主机能否访问 broker 端口、Host/Port/TLS/认证配置是否正确、企业网络或云安全组是否拦截出站 MQTT;若任务运行但查询无数据,检查发布 topic 是否与订阅完全一致、样例能否被识别且字段名与映射匹配、SubTableName 是否设置、目标库与超级表是否创建成功、查询的库名是否正确。

延伸阅读

本主题相关的更多资料可继续阅读当前仓库中的以下文档:

  • No-Code Data Ingestion 目录:支持的数据源、ETL 规则、健康状态与断点续传(本文主体)。
  • Zero-Code Data Ingestion 快速入门:MQTT 端到端零代码接入示例。
  • MQTT 完整配置指南:MQTT 连接、TLS、采集、解析、映射与高级选项的完整说明。
  • 安装 taosX-Agent:在数据源侧网络部署 Agent 代理采集。
  • taosX-Agent 参考:Agent 的部署、配置、运行与排障。
  • Data Ingest and Delivery 章节总览:零代码数据接入、零代码数据投递、数据流管理与迁移指南。
  • 其他数据源专项指南(Kafka、CSV、OPC UA、InfluxDB、OpenTSDB、PI 等)均位于 01-no-code-ingestion 目录 下。 </output文章>

【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询