TDengine 跨集群数据订阅实战:使用 taosExplorer 与 TMQ 实现无代码数据写入
【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine
本指南介绍如何在 TDengine 中通过 taosExplorer(Web 图形界面)与 TMQ(TDengine Message Queue)数据订阅机制,将另一个 TDengine 集群的数据实时同步到当前集群,全程无需编写任何代码。你将掌握从源集群创建 Topic、复制 DSN,到本集群创建"TDengine 数据订阅"任务、配置订阅参数、监控运行状态,以及利用 DSN 多 Topic、自动建 Topic 等高级用法的完整流程。本文以 03-tmq.md 为核心展开,并结合仓库中的 Topic 语法文档与 taosX 组件参考进行深入讲解。
TDengine 的时序数据订阅能力(TMQ)允许应用程序像使用消息队列一样消费数据库中的写入数据。除了通过 编程语言 API 直接订阅,TDengine 企业版还提供了 taosX 组件,它封装了 TMQ 消费逻辑,让你可以在 taosExplorer 中以"零代码"的方式创建数据订阅任务,把源集群的数据持续写入本集群。本文聚焦于"从另一个 TDengine 集群订阅数据到本集群"这一典型场景。
前置条件:taosX 与 XNODE
本文描述的"无代码数据写入"能力由 TDengine 企业版的taosX组件提供(taosX 是 TDengine 企业版的核心组件,提供零代码数据接入能力),相关操作通过 taosExplorer 图形界面完成。在使用前需要确认以下前提:
- 已安装 TDengine 企业版软件包,taosX 以服务模式运行(Linux 下可通过
systemctl start taosx启动,Windows 下可通过sc.exe start taosx启动),具体可参考 taosX 组件参考。 - 从 3.4.0.0 版本开始,taosX 使用 TSDB 作为元数据存储介质(不再使用 SQLite 存储元数据),且创建任务前需要先创建 XNODE(详见 创建 XNODE)。
- 源集群与目标集群之间网络可达,taosX 能够通过 WebSocket 或原生连接访问源集群。
如果 taosX 无法直接连接到数据源所在的网络,也可以先安装 taosX-Agent,再由 taosX 通过 Agent 间接访问数据源。
一、准备工作:在源集群创建 Topic
跨集群订阅的第一步是在源集群上创建订阅所需的 Topic。Topic 是 TMQ 的订阅入口,可以订阅整个数据库(Database)、超级表(Supertable)或子表(Subtable)。本文以订阅一个名为test的数据库为例。
1. 进入"数据订阅"页面
打开源集群的 taosExplorer 界面,点击左侧导航菜单中的"Data Subscription"(数据订阅),进入后点击"Add New Topic"(添加新主题)按钮。
2. 添加新 Topic
在弹出的对话框中输入主题名称(Topic Name),并选择要订阅的数据库。
创建 Topic 时,如果选择订阅的类型为数据库或超级表,且希望同步表的增/删/改(DDL)操作,需要开启同步 Meta(WITH META)选项,用于数据库/超级表的迁移;否则该 Topic 只会进行纯数据同步。
3. 复制 Topic 的 DSN
点击"Create"按钮完成创建,返回主题列表后,复制该 Topic 对应的DSN(Data Source Name,数据源名称)备用。这个 DSN 将在下一步创建订阅任务时粘贴到表单中。
Topic 背后的 SQL 语法
在 taosExplorer 中创建 Topic 只是图形化操作,其底层对应的是 TDengine 的 SQL Topic 语法。理解这些语法有助于你判断"选择哪种订阅对象"以及"是否需要 Meta 同步"。根据 Topic 语法文档,TDengine 从 v3.0.0.0 起支持三类 SQL 创建的 Topic:
- 查询类 Topic(Query Topic):订阅一段 SQL 查询定义的数据流,例如
CREATE TOPIC power_topic AS SELECT ts, current, voltage FROM power.meters WHERE voltage > 200;。查询 Topic 可包含过滤条件和标量函数,但不支持聚合函数、时间窗口聚合以及DISTINCT、GROUP BY、ORDER BY、PARTITION BY、LIMIT/SLIMIT等子句。 - 超级表 Topic(Supertable Topic):订阅指定超级表的全部数据,语法为
CREATE TOPIC [IF NOT EXISTS] topic_name [WITH META | ONLY META] AS STABLE stb_name [where_condition]。其中WITH META会额外返回建超级表及子表的语句(主要用于 taosX 的超级表迁移),ONLY META则只订阅元数据变化而不传输时序数据,where_condition只能用标签或tbname过滤子表。 - 数据库 Topic(Database Topic):订阅指定数据库内所有表的数据,语法为
CREATE TOPIC [IF NOT EXISTS] topic_name [WITH META | ONLY META] AS DATABASE db_name;。WITH META会返回库内所有超级表、子表、普通表的建表/删表/改表语句(主要用于 taosX 的数据库迁移)。
需要注意的是,TDengine 实例中 Topic 的最大数量由配置参数tmqMaxTopicNum控制(默认 20)。如果 Topic 不再需要,可以执行DROP TOPIC [IF EXISTS] [FORCE] topic_name;删除;若有消费者正在订阅,需使用FORCE强制删除。
二、创建订阅任务
Topic 创建完成后,回到目标集群(即本集群),通过 taosExplorer 创建订阅任务。
1. 进入"新增数据源"页面
- 点击左侧菜单"Data Writing"(数据写入)。
- 点击"Add Data Source"(新增数据源)按钮。
2. 输入数据源信息
在表单中完成以下配置:
任务名称(Task Name):输入便于识别的任务名。
任务类型:选择"TDengine Data Subscription"(TDengine 数据订阅)。
目标数据库(Target Database):选择数据要写入的本地数据库。
Topic DSN:将准备阶段复制的 DSN 粘贴到Topic DSN一栏。示例格式为:
tmq+ws://root:taosdata@localhost:6041/topic完成上述步骤后,点击"Connectivity Check"(连通性检查)按钮,测试与源集群的连接是否正常。
关于 DSN 格式的理解
从 taosX 组件参考 可以得知,taosX 使用统一的 DSN(数据源名称)格式描述数据源,其 url-like 结构为:
<driver>[+<protocol>]://[[<username>:<password>]@<host>:<port>][/<object>][?<p1>=<v1>[&<p2>=<v2>]]driver指明数据源驱动,tmq表示通过数据订阅从 TDengine 获取数据;+ws表示通过 WebSocket(REST 接口)获取数据,如果不带+ws则表示使用原生连接,此时 taosx 必须安装在与源集群同网络可达的服务器上;<username>:<password>为源集群的账号密码;<host>:<port>为源集群地址与端口(WebSocket 方式通常为 6041);/<object>为具体的数据源对象,可以是数据库、超级表、子表或 Topic 名称。
因此tmq+ws://root:taosdata@localhost:6041/topic的含义是:通过 WebSocket 连接本地 6041 端口的 TDengine 实例,使用root/taosdata账号,订阅名为topic的 Topic。
3. 填写订阅设置并提交任务
展开订阅选项(Subscribe Options),按需配置以下参数后点击"Submit"(提交):
- 订阅初始位置(Start From):可配置从**最早数据(earliest)或最晚数据(latest)**开始订阅,默认为
earliest。选择earliest意味着 Topic 中从创建(或从头)开始的所有数据都会被同步,适合全量追平场景;选择latest则只订阅任务启动后新写入的数据。 - 超时时间(Timeout):设置订阅超时,支持单位
ms(毫秒)、s(秒)、m(分钟)、h(小时)、d(天)、M(月)、y(年)。 - 订阅组 ID(Group ID):用于标识一个订阅组的任意字符串,最大长度为 192。同一个订阅组内的订阅者共享消费进度(offset);不指定时,taosX 会使用随机生成的 group ID。
- 客户端 ID(Client ID):用于标识客户端的任意字符串,最大长度为 192。
- 同步已落盘数据(TSDB Data):如启用,可以同步已经落盘到 TSDB 时序数据存储文件中(即不在 WAL 中)的数据;如关闭,则只同步尚未落盘(即仍保存在 WAL 中)的数据。
- 同步删表操作(Table Deletions):如启用,则删表操作会被同步到目标数据库。
- 同步删数据操作(Data Deletions):如启用,则数据删除操作会被同步到目标数据库。
- 压缩(Compression):启用 WebSocket 压缩支持,以降低网络带宽占用,适合跨机房或带宽受限的场景。
- 确认无误后点击"Submit"(提交)按钮提交任务。
配置项背后的原理
- 初始位置与超时:TMQ 的消费语义与 Kafka 高度兼容(多语言订阅 API 均与 Kafka 订阅 API 保持高度兼容),
earliest/latest即对应 Kafka 中的auto.offset.reset语义。超时时间则约束了消费端在无新消息时的等待行为。 - 落盘数据同步:这一开关对应 TMQ 对 WAL(Write-Ahead Log)与 TSDB 数据文件的读取能力。TMQ 的消息推送本身基于 WAL;开启"同步已落盘数据"后,taosX 才能读取已经写入 TSDB 文件的历史数据,从而完成全量+增量的一致性同步。
- Meta 同步:开启删表/删数据操作的同步,配合源端 Topic 的
WITH META(或ONLY META)选项,即可实现表结构(DDL)与数据删除操作的跨集群复制,这正是 Topic 语法文档 中描述的数据库/超级表迁移场景。
三、监控任务运行情况
提交任务后,返回"数据源(Data Source)"页面即可查看任务状态。任务会先被加入执行队列,稍后开始运行。
- 点击"View"(查看)按钮,可以监控任务的动态统计信息(Current Metrics),例如
metrics.tmq.records(已消费记录数)、metrics.tmq.points(已消费点数)、metrics.records_per_second(每秒记录数)、metrics.tmq.topics(订阅的 Topic 数)、metrics.tmq.workers(消费工作线程数)等。
- 也可以点击左侧的折叠按钮,展开任务的活动(Activity)信息。如果任务运行异常,这里会给出详细的错误说明,是排查问题(如连通性失败、权限不足、DSN 拼写错误)的第一入口。
可观测性补充:taosX 监控指标
在 taosX 服务模式下,这些任务级指标还会通过 taosKeeper 上报并写入监控数据库。根据 taosX 组件参考,TMQ(TDengine V3 任务)相关的指标包括:
| 指标 | 含义 |
|---|---|
total_messages/messages | TMQ 累计/本次收到的消息总数 |
total_messages_of_meta/messages_of_meta | 收到的 Meta 类型消息数(表结构变化) |
total_messages_of_data/messages_of_data | 收到的 Data 与 MetaData 类型消息数(数据块) |
total_success_blocks/success_blocks | 累计/本次成功写入的数据块数 |
topics | 通过 TMQ 订阅的 Topic 数量 |
consumers | TMQ 消费者数量 |
total_write_raw_fails/write_raw_fails | 原始元数据写入失败次数 |
借助这些指标,你可以在任务运行时快速判断"消费是否跟上写入"以及"是否存在写目标失败"。
四、高级用法
taosExplorer 的"TDengine 数据订阅"数据源还支持以下高级用法,适用于多 Topic、免建 Topic、显式指定消费组等场景:
FROM DSN 支持多个 Topic:多个 Topic 名称用逗号分隔,例如:
tmq+ws://root:taosdata@localhost:6041/topic1,topic2,topic3一个订阅任务即可同时消费多个 Topic 的数据。
在 FROM DSN 中直接使用数据库/超级表/子表名称:可以用数据库名称、超级表名称或子表名称代替 Topic 名称,例如:
tmq+ws://root:taosdata@localhost:6041/db1,db2,db3这种情况下不需要提前创建 Topic,taosX 会自动识别出使用的是数据库名称,并自动在源集群创建订阅对应数据库的 Topic。该能力对"快速打通两个集群、全库搬迁"的场景非常实用。
FROM DSN 支持
group.id参数:可在 DSN 中显式指定订阅所用的 group ID,例如:tmq+ws://root:taosdata@localhost:6041/topic?group.id=my-group不指定时,taosX 会使用随机生成的 group ID。显式指定 group ID 的价值在于:同一 group ID 的多个消费者(或多次运行的任务)可以共享消费进度,便于实现故障恢复后从上次位置续传。
五、从命令行模式理解同源能力
虽然本文聚焦于 taosExplorer 图形化操作,但理解 taosX 的命令行模式有助于你更深刻地认识 DSN 与订阅参数的本质。taosX 命令行格式为:
taosx -f <from-DSN> -t <to-DSN> <other parameters>例如,通过命令行执行一次 TMQ 订阅迁移可以写作:
taosx run -f 'tmq+ws://root:taosdata@localhost:6041/db1' -t 'taos:///db2' -v其中-f指定数据源(Source DSN),-t指定写入目标(Sink DSN),-v将日志级别设为 info(-vv对应 debug,-vvv对应 trace)。--jobs <number>参数可以指定并发任务数(仅支持 tmq 任务)。在图形界面中创建的任务,本质上就是 taosX 服务模式托管执行的同类任务,二者共用同一套 DSN 语义与订阅机制。
总结
跨集群数据订阅是 TDengine 多集群架构下"数据汇聚与容灾"的常用手段。通过本文你可以看到,整个过程完全在 taosExplorer 界面内完成:
- 源集群:创建 Topic(数据库/超级表/子表均可,按需开启 Meta 同步)→ 复制 DSN;
- 本集群:新增"TDengine 数据订阅"数据源 → 粘贴 DSN → 配置初始位置、超时、group ID、client ID、落盘/删表/删数据同步与压缩等选项 → 提交任务;
- 监控:通过任务列表、Current Metrics 与活动信息持续观测,必要时借助 taosKeeper 的 TMQ 指标深挖;
- 进阶:多 Topic 逗号分隔、直接用库名/表名免建 Topic、DSN 显式指定
group.id。
如需进一步了解 Topic 的 SQL 语法细节(查询/超级表/数据库三类 Topic、WITH META/ONLY META、RELOAD TOPIC、消费者组管理等),可继续阅读 Topic 语法文档;taosX 服务配置(taosx.toml)、命令行参数与完整监控指标见 taosX 组件参考;多语言订阅 API 的编程方式见 开发者指南·数据订阅。
【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考