TDengine 跨集群数据订阅实战:基于 TMQ 与 taosExplorer 实现源集群到本集群的零代码数据同步
TDengine 跨集群数据订阅实战基于 TMQ 与 taosExplorer 实现源集群到本集群的零代码数据同步【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine本文介绍如何使用 taosExplorer 的「TDengine 数据订阅」数据源通过 TDengine 内置的消息队列 TMQ 将另一个集群源集群的 Topic 数据持续订阅写入本集群目标集群。读完后你将掌握完整的跨集群同步操作链路在源集群创建 Topic含 Meta 同步选项、通过 Topic DSN 配置订阅任务、理解订阅初始位置/订阅组/落盘数据同步等关键参数并能在 taosExplorer 中监控任务运行状态与排错。方案背景与适用前提TDengine 从v3.0.0.0开始对消息队列做了大幅优化和增强用户可以在源集群通过 SQL 或 taosExplorer 创建订阅主题Topic再由目标集群的 taosX 组件作为消费者拉取 Topic 数据并写入本集群。TMQ 的订阅 API 与 Kafka 订阅 API 保持高度一致便于复用既有开发经验。本方案属于“零代码接入”全程通过 taosExplorer 图形界面完成无需编写任何同步代码。适用前提与限制taosExplorer 的「数据写入」图形能力由 taosX 组件提供服务模式任务需要先安装 TDengine TSDB Enterprise 安装包并启动 taosX 服务详细说明见 taosX 参考手册源集群与目标集群之间需保证网络连通目标集群通过 WebSocket REST 端口默认 6041访问源集群的 TMQ一个 TDengine 实例可创建的 Topic 个数上限由tmqMaxTopicNum参数控制默认值为 20详见 taosd 配置参数。整体数据流向为源集群 Topic数据库/超级表/子表数据 可选 Meta 语句→ TMQ 消息队列 → 目标集群 taosX 订阅任务 → 目标集群指定数据库。准备工作在源集群创建订阅 Topic在源集群创建订阅所需的 Topic可以订阅整个库、超级表或子表。本示例中我们演示订阅一个名为test的数据库。第一步进入“数据订阅”页面打开源集群的 taosExplorer 界面点击左侧“数据订阅”菜单然后点击“添加新主题”。第二步添加新主题输入主题名称选择要订阅的数据库。对于数据库或超级表类型如果需要同步表的增/删/改操作则需要开启同步 Meta选项用于数据库/超级表的迁移否则此主题只会进行数据同步。这一点与 Topic 的类型定义直接相关数据库主题订阅一个数据库里所有数据等价 SQL 为CREATE TOPIC [IF NOT EXISTS] topic_name [WITH META | ONLY META] AS DATABASE db_name;其中WITH META参数指定后将返回数据库里所有超级表、子表、普通表的元数据创建、删除、修改语句主要用于 taosX 做数据库迁移超级表主题订阅一个超级表中的所有数据WITH META指定后将返回创建超级表、子表等语句主要用于 taosX 做超级表迁移查询主题订阅一条 SQL 查询定义的数据流一旦创建订阅结构即确定。关于主题的完整创建语法、删除与查看、消费者与回放等详情请参考 主题语法。说明超级表订阅和库订阅属于高级订阅模式容易出错如确实要使用请咨询技术支持人员。第三步复制主题的 DSN点击“创建”按钮回到主题列表复制主题的DSN备用。DSNData Source Name采用类 URL 格式driver[protocol]://[[username:password]host:port][/object][?p1v1[p2v2]]对于 TMQ 订阅源driver为tmqws表示使用 WebSocket 方式获取数据不使用ws则表示使用原生连接此时 taosX 所在服务器需要安装 taosc。创建订阅任务在目标集群配置“TDengine 数据订阅”数据源第一步进入“新增数据源”页面点击左侧“数据写入”菜单点击“新增数据源”第二步输入数据源信息输入任务名称选择任务类型“TDengine 数据订阅”选择目标数据库粘贴准备步骤复制的 DSN 到Topic DSN一栏。例如tmqws://root:taosdatalocalhost:6041/topic完成以上步骤点击“连通性检查”按钮测试与源端的连通性第三步填写订阅设置并提交任务选择订阅初始位置。可配置从最早数据earliest或最晚latest数据开始订阅默认为earliest设置超时时间。支持单位 ms毫秒、s秒、m分钟、h小时、d天、M月、y年设置订阅组 ID。订阅组 ID 是用于标识一个订阅组的任意字符串最大长度为 192。同一个订阅组内的订阅者共享消费进度。不指定情况下将使用随机生成的 group ID设置客户端 ID。客户端 ID 是一个用于标识客户端的任意字符串最大长度为 192同步已落盘数据。如启用可以同步已经落盘到 TSDB 时序数据存储文件中即不在 WAL 中的数据。如关闭则只同步尚未落盘即保存在 WAL 中的数据同步删表操作。如启用则会同步删表操作到目标数据库同步删数据操作。如启用则会同步删数据操作到目标数据库压缩。启用 WebSocket 压缩支持以降低网络带宽占用点击“提交”按钮提交任务这些图形界面上的选项本质上就是 TMQ 订阅 DSN 的查询参数理解参数对应关系有助于排查问题或在 taosX 命令行模式下复用相同配置界面选项对应 DSN 参数说明订阅初始位置auto.offset.reset取earliest或latest默认earliest订阅组 IDgroup.id订阅的组 ID同组内订阅者共享消费进度不指定则随机生成客户端 IDclient.id标识客户端的任意字符串选填同步已落盘数据experimental.snapshot.enable启用后可同步已落盘到 TSDB 时序数据存储文件不在 WAL 中的数据关闭则只同步保存在 WAL 中的未落盘数据同步删表操作with.meta.delete/with.meta同步元数据中的删除表事件仅当元数据同步启用时有效同步删数据操作with.meta.drop/with.meta同步元数据中的删除数据事件仅当元数据同步启用时有效压缩WebSocket 压缩降低网络带宽占用参数解析的实现可以在客户端 TMQ 配置模块中找到clientTmqConf.c 中对group.id、client.id、auto.offset.reset、experimental.snapshot.enable等键值逐一校验并赋值例如group.id不允许包含:字符auto.offset.reset仅接受合法取值。这也解释了界面上“最大长度 192”等约束的来源。需要注意语义差异with.meta.*系列参数对应的是 Topic 侧WITH META产生的 Meta 事件建表/改表/删表/删数据等元数据语句在消费端的落地开关而 Topic 创建时的同步 Meta选项决定了源端是否把元数据语句放进消息流。两端配合才能完成数据库/超级表结构的迁移同步。监控任务运行情况提交任务后回到数据源页面可以查看任务状态。任务先会被加入执行队列稍后就开始运行。点击“查看”按钮可以监控任务的动态统计信息。也可以点击左侧折叠按钮展开任务的活动信息。如果任务运行异常此处可以看到详细的说明。这些动态统计信息对应 taosX 上报给 taosKeeper 的「taosX TDengine V3 任务」监控指标例如total_messages通过 TMQ 累计收到的消息总数、total_messages_of_metaMeta 类型消息数、total_messages_of_dataData 和 MetaData 类型消息数、total_write_raw_fails写入 raw meta 失败的次数、topics订阅的主题数、consumersTMQ 消费者数等完整字段说明见 taosX 参考手册 · 监控指标。此外taosExplorer 数据源页面中“查看”的动态信息也便于直观确认订阅任务是否持续消费到数据。高级用法FROM DSN 支持多个 Topic多个 Topic 的名字用逗号分割。例如tmqws://root:taosdatalocalhost:6041/topic1,topic2,topic3在 FROM DSN 中也可以用数据库名称、超级表名称或子表名称代替 Topic 名称。例如tmqws://root:taosdatalocalhost:6041/db1,db2,db3此时不必提前创建 TopictaosX 将自动识别到使用的是数据库名称并自动在源集群创建订阅数据库的 Topic。FROM DSN 支持group.id参数以显式指定订阅用的 group ID。不指定情况下将使用随机生成的 group ID。显式指定group.id的实践意义在于消费进度管理同一订阅组内的订阅者共享消费进度因此用固定的 group ID 可以在任务重启后从上次进度继续消费避免使用随机 group ID 导致每次重启都从初始位置默认earliest重新拉取全量数据。底层实现与延伸阅读从源码结构看TMQ 的订阅/投递能力分布在几处关键模块客户端配置与消费者接口clientTmqConf.c 负责订阅参数解析多语言连接器 API 的订阅用法见 开发指南 · 数据订阅服务端 MQTT 桥接库source/libs/tmqtt 目录包含 TMQ 消息经 MQTT 协议发布的实现与示例其中 tmqttMgmt.c 展示了 taosd 侧管理taosmqtt子进程的启动逻辑官方示例 sub.py 演示了使用 paho-mqtt 通过共享订阅$share/group/topic消费 TMQ Topic配套的建库建 Topic 脚本为 prep.sql。从v3.3.7.0开始还提供 MQTT 订阅功能可通过 MQTT 客户端直接订阅数据详见 MQTT 订阅taosX 命令行模式下同样的 TMQ 订阅能力可写成taosx run -f tmqws://username:passwordip:port/topic?paramvalue... -t taosws://...TMQ DSN 全部参数group.id、client.id、auto.offset.reset、with.meta、experimental.snapshot.enable等见 taosX 参考手册命令行还额外支持--jobs指定 tmq 任务的并发数若不需要图形界面也可以直接在taosshell 中执行subscribe topic -g group_id创建消费者快速验证订阅功能用法见 taos CLI 数据订阅。综合来看本篇覆盖的「taosExplorer 图形化 TMQ 数据订阅」方案适合集群间持续数据复制与异地容灾场景Topic 侧决定订阅范围库/超级表/子表 是否携带 Meta订阅任务侧决定消费行为初始位置、消费组、落盘数据回放、Meta 事件落地与网络压缩两者通过 DSN 串联配合 taosX 的监控指标即可完成从配置到运维的闭环。【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考