ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

TDengine 基于 MQTT 的数据订阅:Bnode 管理与 taosmqtt 消费实践

2026/9/12 17:22:06 拓冰建站 浏览量
TDengine 基于 MQTT 的数据订阅:Bnode 管理与 taosmqtt 消费实践 TDengine 基于 MQTT 的数据订阅Bnode 管理与 taosmqtt 消费实践【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine自v3.3.7.0起TDengine 原生支持通过 MQTT 协议进行数据订阅使用任意兼容 MQTT 的客户端连接到 TDengine 的 Bnode 服务即可直接订阅已存在的 Topic 数据。本文将完整介绍 Bnode 的创建、查询与删除操作以 Pythonpaho-mqtt为例演示订阅全流程并结合仓库源码bnode.c、tmqttMgmt.c 等解析其工作原理、消息格式与关键配置参数帮助你快速在工业物联网IIoT场景下搭建基于 MQTT 的实时数据消费链路。功能特性概览TDengine 的 MQTT 订阅具备以下核心能力协议支持推荐使用 MQTT 5.0同时兼容 MQTT 3.1 与 3.1.1。需要注意的是sub-offset等用户属性User Property依赖 MQTT 5.0低版本协议无法使用。认证方式复用 TDengine 原生认证体系使用数据库账号密码即可连接无需额外维护认证服务。Topic 管理与标准 MQTT 协议不同TDengine 的 Topic 必须预先创建由于不支持消息发布Topic 无法通过消息发布动态创建创建方式见后文 SQL 示例。共享订阅形如$share/group_id/topic_name的 Topic 会被识别为共享订阅适用于需要负载均衡与高可用的消费场景。订阅位置支持latest默认与earliest最早 WAL 位置两种起始位置。通过订阅用户属性sub-offsetearliest请求最早位置。服务质量支持 QoS 0 与 QoS 1。Bnode 管理BnodeBroker Node是 TDengine 集群中负责提供 MQTT 订阅服务的节点组件可通过taosCLI 进行管理。创建 Bnode使用如下 SQL 语句创建 BnodeCREATE BNODE ON DNODE dnode_id;每个 dnode 上只能创建一个 bnode。Bnode 创建成功后会自动启动名为taosmqtt的 bnode 子进程用于提供 MQTT 订阅服务。从源码看这一过程由 bnode.c 中的bndOpen()完成当协议为TSDB_BNODE_OPT_PROTO_MQTT时调用mqttMgmtStartMqttd()该函数通过 libuv 的进程管理能力uv_spawn相关逻辑见 tmqttMgmt.c拉起taosmqtt可执行文件Linux 下位于/usr/bin/taosmqttWindows 下为C:\TDengine\taosmqtt.exe。taosmqtt服务默认使用6057端口。如需更换端口可修改taos.cfg中的mqttPort参数。该参数在 tglobal.c 中注册为整型配置项取值范围1 ~ 65056作用域为服务端CFG_SCOPE_SERVER即需要修改 taosd 配置文件并重启后生效。查看 Bnode使用如下 SQL 语句查看集群中的 Bnode 信息完整字段列表参见INS_BNODESSHOW BNODES;输出类似taos show bnodes; id | endpoint | protocol | create_time | 1 | 192.168.0.1:6057 | mqtt | 2024-11-28 18:44:27.089 | Query OK, 1 row(s) in set (0.037205s)可以看到每个 Bnode 包含id、endpointdnode 地址与 mqttPort 端口、protocol此处为mqtt以及create_time等字段。删除 Bnode使用如下 SQL 语句删除 BnodeDROP BNODE ON DNODE dnode_id;删除操作会将该 bnode 从 TDengine 集群中移除同时停止对应的taosmqtt服务。对应源码中的bndClose()会调用mqttMgmtStopMqttd()停止 MQTT 守护进程见 bnode.c。MQTT 数据订阅示例下面通过一个完整示例演示先在 TDengine 中创建测试数据再订阅这些数据。示例使用 MQTT 5.0 与 Pythonpaho-mqtt库以便设置sub-offset用户属性你也可以使用任何兼容 MQTT 的客户端。仓库自带的完整可运行示例位于 source/libs/tmqtt/example 目录含prep.sql、sub.py、run.sh可直接参考。创建测试数据在 taos CLI 中执行以下 SQL 语句完成示例环境的准备CREATE DATABASE db VGROUPS 1; CREATE TABLE db.meters (ts TIMESTAMP, f1 INT) TAGS (t1 INT); CREATE TOPIC topic_meters AS SELECT ts, tbname, f1, t1 FROM db.meters; INSERT INTO db.tb USING db.meters TAGS (1) VALUES (now, 1); CREATE BNODE ON DNODE 1;以上语句依次完成了创建数据库1 个 vgroup、创建超级表、基于超级表创建 TopicCREATE TOPIC语句将数据映射为可订阅的消息流、插入一条测试数据最后在 dnode 1 上创建 Bnode 启动订阅服务。编写消费者将以下 Python 代码保存为sub.pyimport time import paho.mqtt import paho.mqtt.properties as p import paho.mqtt.packettypes as pt import paho.mqtt.client as mqttClient def on_connect(client, userdata, flags, rc, propertiesNone): print(CONNACK received with code %s. % rc) sub_properties p.Properties(pt.PacketTypes.SUBSCRIBE) sub_properties.UserProperty (sub-offset, earliest) client.subscribe($share/g1/topic_meters, qos1, propertiessub_properties) def on_subscribe(client, userdata, mid, granted_qos, propertiesNone): print(Subscribed: str(mid) str(granted_qos)) def on_message(client, userdata, msg): print(msg.topic str(msg.qos) str(msg.payload)) if paho.mqtt.__version__[0] 1: client mqttClient.Client(mqttClient.CallbackAPIVersion.VERSION2, client_idtmq_sub_cid, userdataNone, protocolmqttClient.MQTTv5) else: client mqttClient.Client(client_idtmq_sub_cid, userdataNone, protocolmqttClient.MQTTv5) client.on_connect on_connect client.username_pw_set(root, taosdata) client.connect(127.0.1.1, 6057) client.on_subscribe on_subscribe client.on_message on_message client.loop_forever()关键点说明连接信息使用 TDengine 原生账号认证root/taosdata连接到本机6057端口即taosmqtt服务地址可按实际 dnode 地址调整。订阅 Topic$share/g1/topic_meters是共享订阅形式g1为消费组名topic_meters为预先创建的 Topic。订阅位置通过 SUBSCRIBE 包的用户属性设置(sub-offset, earliest)请求从最早 WAL 位置开始消费不设置时默认从latest开始。这一属性在服务端最终映射为 TDengine 消费者TMQ的auto_offset_reset配置earliest/latest见 tmqttCtx.c。QoS示例使用 QoS 1至少一次投递。运行订阅依次执行以下命令安装依赖并启动消费者python3 -m venv .test-env source .test-env/bin/activate pip3 install paho-mqtt2.1.0 python3 ./sub.py订阅成功后之后写入topic_meters的任何新数据都会被自动推送到客户端。消息格式以上一节示例为例客户端会输出类似信息CONNACK received with code Success. Subscribed: 1 [ReasonCode(Suback, Granted QoS 1)] topic_meters 1 b{topic:topic_meters,db:db,vid:2,rows:[{ts:1753086482326,tbname:tb,f1:1,t1:1}]}对第三行逐段解读topic_meters本次订阅的 Topic 名称1该消息的 QoS 值之后是 UTF-8 编码的 JSON 消息体各字段含义如下字段含义topic消息所属 Topic 名称db数据所在的数据库名vid数据所在的 vgroup虚拟节点组IDrows数据行数组每行包含建 Topic 时 SELECT 出的全部列如ts毫秒时间戳、tbname子表名、f1、t1等进阶集群化负载均衡与测试工具共享订阅$share/group_id/topic在多个消费者同属一个消费组时消息会在组内分发从而在多个消费客户端之间实现负载均衡与高可用不同消费组之间则各自独立消费同一份数据。如果需要压测或快速验证订阅链路仓库 source/libs/tmqtt/tools 下提供了相关工具topic-producer.c用于向 TDengine Topic 写入测试数据的生产者程序其中同样配置了auto_offset_reset为earliest/latest等消费参数perf.py基于paho-mqtt的订阅压测脚本演示了$share/g3/topic_meters等不同共享消费组的订阅写法。小结通过 Bnode 组件TDengine 将内置的 Topic 数据流以标准 MQTT 协议对外暴露使任意 MQTT 生态客户端都能直接消费时序数据大幅降低了与消息中间件集成的门槛。实际部署时需注意Topic 必须预先通过 SQL 创建追求sub-offset等高级特性请使用 MQTT 5.0mqttPort端口可在taos.cfg中按需调整。更多 taosd 配置说明参见 taosd 配置参数Bnode 元数据字段参见INS_BNODES。【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考