- 数据库
- 时序数据库
- 物联网
- 大数据
- 实时分析
- 云原生
【免费下载链接】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.
本篇技术指南讲解如何在 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 BY、ORDER BY、LIMIT等子句;订阅数据结构在创建时固定。 - 超级表主题(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) |
topic | object 部分,即主题名称 |
三、创建订阅任务
Topic 就绪后,回到本集群的 taosExplorer 创建订阅任务,将源集群的数据写入本集群。
第一步:进入「新增数据源」页面
- 点击左侧「数据写入」菜单;
- 点击「新增数据源」。
第二步:输入数据源信息
- 输入任务名称;
- 选择任务类型为「TDengine 数据订阅」;
- 选择目标数据库(数据将写入本集群的哪个库);
- 将准备阶段复制的 DSN 粘贴到Topic DSN一栏,例如:
tmq+ws://root:taosdata@localhost:6041/topic; - 完成以上步骤后,点击「连通性检查」按钮,测试与源端的连通性,确保账号权限、网络与端口均可达。
第三步:填写订阅设置并提交任务
订阅设置项是控制同步行为的关键,逐项说明如下:
- 订阅初始位置(offset reset):可配置从最早数据(
earliest)或最晚(latest)数据开始订阅,默认为earliest。选择earliest会将源端 Topic 中尚未消费的历史数据一并同步;选择latest则仅同步任务启动之后产生的新数据。 - 超时时间:设置订阅/消费相关的超时时长,支持单位
ms(毫秒)、s(秒)、m(分钟)、h(小时)、d(天)、M(月)、y(年)。 - 订阅组 ID(group.id):用于标识一个订阅组的任意字符串,最大长度为 192。同一个订阅组内的订阅者共享消费进度——这正是「断点续传」的机制:任务重启后从该组已提交的 offset 继续消费,避免重复或遗漏。不指定时,将使用随机生成的 group ID。
- 客户端 ID(client.id):用于标识客户端的任意字符串,最大长度为 192,便于在多客户端场景下区分连接。
- 同步已落盘数据:如启用,可以同步已经落盘到TSDB时序数据存储文件中(即不在 WAL 中)的数据;如关闭,则只同步尚未落盘(即保存在WAL中)的数据。TDengine 的写入路径是先写 WAL(Write-Ahead Log)再落盘到 TSDB 文件,因此这一选项决定了同步范围是否覆盖历史已落盘数据。
- 同步删表操作:如启用,则会同步删表操作到目标数据库,保持两端表结构一致。
- 同步删数据操作:如启用,则会同步删数据操作到目标数据库,保证目标端数据删除语义与源端一致。
- 压缩:启用 WebSocket 压缩支持,以降低跨集群同步时的网络带宽占用,适合数据量大、带宽受限的场景。
- 点击「提交」按钮,提交任务。
四、监控任务运行情况
提交任务后,回到「数据写入 → 数据源」页面即可查看任务状态。任务会先被加入执行队列,稍后开始运行。
点击「查看」按钮,可以监控任务的动态统计信息。这些指标由 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 还支持以下高级用法,适用于批量迁移、多对象同步等场景:
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-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.id、auto.offset.reset、enable.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.
相关推荐
TDengine 跨集群数据订阅实战:基于 TMQ 与 taosExplorer 实现源集群到本集群的零代码数据同步
TDengine 跨集群数据订阅实战:基于 TMQ 与 taosExplorer 实现源集群到本集群的零代码数据同步 本文介绍如何使用 taosExplorer
数据库时序数据库大数据物联网云原生TDengine 跨集群数据订阅实战:使用 taosExplorer 与 TMQ 实现无代码数据写入
TDengine 跨集群数据订阅实战:使用 taosExplorer 与 TMQ 实现无代码数据写入 本指南介绍如何在 TDengine 中通过 taosExp
数据库时序数据库大数据物联网云原生TDengine 可视化管理实战:taosExplorer 零代码管理集群、数据、流计算与订阅
TDengine 可视化管理实战:taosExplorer 零代码管理集群、数据、流计算与订阅 taosExplorer 是 TDengine v3.0 起随服
数据库时序数据库大数据物联网云原生
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考