TDengine 零代码数据订阅实战:使用 taosExplorer 与 TMQ 将集群数据同步到本集群
2026/9/22 0:56:12 网站建设 项目流程
  • 数据库
  • 时序数据库
  • 物联网
  • 大数据
  • 实时分析
  • 云原生

【免费下载链接】tdengine

TDengine is an open source, high-performance, cloud native time-series database optimized for Internet of Things (IoT), Connected Cars, Industrial IoT and DevOps.

项目地址:https://gitcode.com/taosdata/tdengine
点击查看免费下载

本篇技术指南讲解如何在 TDengine 中使用taosExplorer + TMQ(TDengine Message Queue)以零代码方式创建数据订阅数据源,将另一个集群的数据实时同步到本集群。你将从零开始完成「在源集群创建 Topic → 复制 DSN → 在本集群创建订阅任务 → 监控运行」的完整闭环,并深入理解 Topic 类型、WAL/TSDB 数据边界、订阅组与消费进度等底层原理,从而能独立完成跨集群数据订阅与迁移的落地实施。

说明:本文描述的「TDengine 数据订阅」数据源由 TDengine Enterprise 的taosX组件提供,功能入口集成在 taosExplorer 图形界面中,其服务模式与 DSN 规范可参考 taosX 参考文档。

一、理解 TMQ 数据订阅与零代码接入

TDengine 从v3.0.0.0起对消息队列能力进行了大幅优化与增强,数据订阅(TMQ)是其内置能力:Topic由 SQL 创建,连接器(C、Java、Go、Rust、Python、C# 等)、taos命令行乃至 MQTT 客户端都可以像消费消息队列一样订阅并消费数据,消费进度以「订阅组(consumer group)」为单位推进。完整概念可参见 数据订阅主题语法 与 原生订阅。

在此基础上,零代码数据接入把「订阅」与「写入」两个环节封装为可视化任务:

  • 源端:在源集群通过 SQL 或 taosExplorer 创建 Topic,Topic 可以订阅整个数据库、超级表或子表;
  • 目标端:在本集群的 taosExplorer「数据写入」中新增一个TDengine 数据订阅数据源,填入源端 Topic 的 DSN,即可由 taosX 自动完成跨集群的数据订阅、解析与写入;
  • 全程无需编写任何代码,任务状态、动态统计信息与异常活动详情都可以在界面中直接查看。

从源码实现看,Topic 的创建最终是一条走向 mnode 的命令:客户端 SQL 经 parser 层解析并生成QUERY_NODE_CREATE_TOPIC_STMT语法树,再经 parTranslater 封装为TDMT_MND_TMQ_CREATE_TOPIC命令 下发,最终由 mnode 的mndProcessCreateTopicReq完成鉴权、校验并持久化 Topic 元数据。本文的界面操作本质上是这一底层链路的可视化封装。

二、准备工作:在源集群创建订阅 Topic

创建订阅任务之前,需要先在源集群准备好可供订阅的 Topic。Topic 可订阅整个数据库、超级表或子表,本文示例演示订阅一个名为test的数据库。

第一步:进入「数据订阅」页面

打开源集群的 taosExplorer 界面,点击左侧「数据订阅」菜单,然后点击「添加新主题」。

第二步:添加新主题

输入主题名称,选择要订阅的数据库(也可选择超级表或子表粒度)。

对于数据库超级表类型的 Topic,如果需要同步表的增/删/改操作,需要开启同步 Meta选项,用于数据库/超级表的整体迁移;否则该 Topic 只会进行数据同步,不会同步表结构的变更。

关于 Topic 创建的完整 SQL 语法(三种 Topic 类型、WITH META/ONLY META语义、删除/查看/重载 Topic 等),请参考 主题语法,这里简要说明与迁移场景最相关的三类 Topic:

  • 查询主题(Query Topic)CREATE TOPIC ... AS subquery,订阅一条 SQL 查询定义的数据流,可带过滤条件与标量函数,但不支持聚合、窗口、GROUP BYORDER BYLIMIT等子句;订阅数据结构在创建时固定。
  • 超级表主题(Supertable Topic)CREATE TOPIC ... AS STABLE stb_name,订阅指定超级表的全部数据;WITH META会额外返回建超级表/建子表的语句(taosX 迁移超级表主要依赖此模式),ONLY META则只订阅元数据变更,不传输时序数据。
  • 数据库主题(Database Topic)CREATE TOPIC ... AS DATABASE db_name,订阅数据库内所有表的数据;同样支持WITH META/ONLY META,是 taosX 做数据库迁移(含库内所有超级表、子表、普通表的元数据变更)的主要模式。

第三步:复制主题的 DSN

点击「创建」按钮完成创建,回到主题列表,复制该主题的DSN备用。DSN 是后续创建订阅任务时连接源端的关键信息,其格式为:

<driver>[+<protocol>]://[[<username>:<password>@]<host>:<port>][/<object>][?<p1>=<v1>[&<p2>=<v2>]]

tmq驱动为例,一个典型的 Topic DSN 形如:

tmq+ws://root:taosdata@localhost:6041/topic

各组成部分的含义(与 taosX 参考文档 中的 DSN 规范一致):

组成部分含义
tmq驱动(driver),表示使用数据订阅(TMQ)方式从 TDengine 取数;若使用taos则表示通过查询接口取数
+ws协议(protocol),表示通过 REST/WebSocket 连接;不带+ws则使用原生连接,此时 taosx 必须部署在服务器上
root:taosdata源集群的用户名与密码
localhost:6041源集群地址与端口(6041 为 WebSocket/REST 端口;原生连接通常为 6030)
topicobject 部分,即主题名称

三、创建订阅任务

Topic 就绪后,回到本集群的 taosExplorer 创建订阅任务,将源集群的数据写入本集群。

第一步:进入「新增数据源」页面

  1. 点击左侧「数据写入」菜单;
  2. 点击「新增数据源」。

第二步:输入数据源信息

  1. 输入任务名称
  2. 选择任务类型为「TDengine 数据订阅」;
  3. 选择目标数据库(数据将写入本集群的哪个库);
  4. 将准备阶段复制的 DSN 粘贴到Topic DSN一栏,例如:tmq+ws://root:taosdata@localhost:6041/topic
  5. 完成以上步骤后,点击「连通性检查」按钮,测试与源端的连通性,确保账号权限、网络与端口均可达。

第三步:填写订阅设置并提交任务

订阅设置项是控制同步行为的关键,逐项说明如下:

  1. 订阅初始位置(offset reset):可配置从最早数据(earliest)或最晚(latest)数据开始订阅,默认为earliest。选择earliest会将源端 Topic 中尚未消费的历史数据一并同步;选择latest则仅同步任务启动之后产生的新数据。
  2. 超时时间:设置订阅/消费相关的超时时长,支持单位ms(毫秒)、s(秒)、m(分钟)、h(小时)、d(天)、M(月)、y(年)。
  3. 订阅组 ID(group.id):用于标识一个订阅组的任意字符串,最大长度为 192。同一个订阅组内的订阅者共享消费进度——这正是「断点续传」的机制:任务重启后从该组已提交的 offset 继续消费,避免重复或遗漏。不指定时,将使用随机生成的 group ID。
  4. 客户端 ID(client.id):用于标识客户端的任意字符串,最大长度为 192,便于在多客户端场景下区分连接。
  5. 同步已落盘数据:如启用,可以同步已经落盘到TSDB时序数据存储文件中(即不在 WAL 中)的数据;如关闭,则只同步尚未落盘(即保存在WAL中)的数据。TDengine 的写入路径是先写 WAL(Write-Ahead Log)再落盘到 TSDB 文件,因此这一选项决定了同步范围是否覆盖历史已落盘数据。
  6. 同步删表操作:如启用,则会同步删表操作到目标数据库,保持两端表结构一致。
  7. 同步删数据操作:如启用,则会同步删数据操作到目标数据库,保证目标端数据删除语义与源端一致。
  8. 压缩:启用 WebSocket 压缩支持,以降低跨集群同步时的网络带宽占用,适合数据量大、带宽受限的场景。
  9. 点击「提交」按钮,提交任务。

四、监控任务运行情况

提交任务后,回到「数据写入 → 数据源」页面即可查看任务状态。任务会先被加入执行队列,稍后开始运行。

点击「查看」按钮,可以监控任务的动态统计信息。这些指标由 taosX 周期上报给 taosKeeper,写入监控库(默认log库)中。对于 TMQ 类任务,重点关注的指标包括(字段定义见 taosX 参考文档 · 监控指标):

指标含义
total_messages/messages累计 / 本次运行通过 TMQ 接收的消息总数
total_messages_of_data/messages_of_data累计 / 本次运行的 Data 与 MetaData 类型消息数
total_success_blocks/success_blocks累计 / 本次运行成功写入的数据块数
topics/consumers通过 TMQ 订阅的 Topic 数与消费者数
total_write_raw_fails/write_raw_fails累计 / 本次运行写入 raw meta 的失败次数

也可以点击左侧折叠按钮,展开任务的活动信息。如果任务运行异常,此处可以看到详细的错误说明,是排查问题(如连通性失败、权限不足、目标库不存在等)的第一现场。

五、高级用法

除界面化配置外,TMQ 数据订阅的 DSN 还支持以下高级用法,适用于批量迁移、多对象同步等场景:

  1. FROM DSN 支持多个 Topic:多个 Topic 名称用逗号分隔,一次任务订阅多个主题。例如:

    tmq+ws://root:taosdata@localhost:6041/topic1,topic2,topic3
  2. 用库名/超级表名/子表名代替 Topic 名称:在 FROM DSN 中,可以直接使用数据库名称、超级表名称或子表名称代替 Topic 名称。例如:

    tmq+ws://root:taosdata@localhost:6041/db1,db2,db3

    此时不必提前在源集群创建 Topic,taosX 会自动识别出使用的是数据库名称,并自动在源集群创建订阅该数据库的 Topic,显著简化了批量数据库迁移的准备工作。

  3. FROM DSN 支持group.id参数:可在 DSN 中显式指定订阅用的 group ID,例如:

    tmq+ws://root:taosdata@localhost:6041/topic?group.id=my-sub-group

    不指定时,将使用随机生成的 group ID。显式指定后,任务可复用既有订阅组的消费进度(offset),实现任务重建后的无缝续传。

六、小结与延伸阅读

本文完整覆盖了「源集群建 Topic → 复制 DSN → 本集群建订阅任务 → 监控与排障」的零代码跨集群数据订阅流程,并解析了订阅位置、超时、订阅组、WAL/TSDB 数据边界、删表/删数据同步、WebSocket 压缩等关键选项的语义。需要强调的是,订阅组(group.id)是消费进度的载体,DSN 中的group.id参数与界面中的「订阅组 ID」一脉相承,理解它是用好断点续传与多 Topic 批量迁移的前提。

如需进一步深入,建议结合以下仓库文档阅读:

  • 主题语法(Topic Syntax):三种 Topic 类型的 SQL 语法、WITH META/ONLY META语义、RELOAD TOPIC等;
  • 原生订阅(Native Subscription):消费者创建、poll、offset 提交等消费模型与关键参数;
  • 数据订阅 API(Managing Consumers):各语言连接器的消费者参数表(group.idauto.offset.resetenable.auto.commit等);
  • taosX 参考文档:DSN 规范、服务模式部署(taosx.toml配置)、命令行模式与监控指标;
  • 源码实现参考:Topic 语法树构造、Topic 命令封装、mnode 创建 Topic 处理。
  • 数据库
  • 时序数据库
  • 物联网
  • 大数据
  • 实时分析
  • 云原生

【免费下载链接】tdengine

TDengine is an open source, high-performance, cloud native time-series database optimized for Internet of Things (IoT), Connected Cars, Industrial IoT and DevOps.

项目地址:https://gitcode.com/taosdata/tdengine
点击查看免费下载

相关推荐

上一篇:TimeMixer项目中的多特征时间序列预测实现方法
下一篇:LiteLoaderQQNT-Anti-Recall插件页面空白问题解决方案

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

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

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

立即咨询