news 2026/9/12 17:21:59

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
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 订阅具备以下核心能力:

  1. 协议支持:推荐使用 MQTT 5.0;同时兼容 MQTT 3.1 与 3.1.1。需要注意的是,sub-offset等用户属性(User Property)依赖 MQTT 5.0,低版本协议无法使用。
  2. 认证方式:复用 TDengine 原生认证体系,使用数据库账号密码即可连接,无需额外维护认证服务。
  3. Topic 管理:与标准 MQTT 协议不同,TDengine 的 Topic 必须预先创建(由于不支持消息发布,Topic 无法通过消息发布动态创建),创建方式见后文 SQL 示例。
  4. 共享订阅:形如$share/group_id/topic_name的 Topic 会被识别为共享订阅,适用于需要负载均衡与高可用的消费场景。
  5. 订阅位置:支持latest(默认)与earliest(最早 WAL 位置)两种起始位置。通过订阅用户属性sub-offset=earliest请求最早位置。
  6. 服务质量:支持 QoS 0 与 QoS 1。

Bnode 管理

Bnode(Broker Node)是 TDengine 集群中负责提供 MQTT 订阅服务的节点组件,可通过taosCLI 进行管理。

创建 Bnode

使用如下 SQL 语句创建 Bnode:

CREATE 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/taosmqtt,Windows 下为C:\TDengine\taosmqtt.exe)。

taosmqtt服务默认使用6057端口。如需更换端口,可修改taos.cfg中的mqttPort参数。该参数在 tglobal.c 中注册:为整型配置项,取值范围1 ~ 65056,作用域为服务端(CFG_SCOPE_SERVER),即需要修改 taosd 配置文件并重启后生效。

查看 Bnode

使用如下 SQL 语句查看集群中的 Bnode 信息(完整字段列表参见INS_BNODES):

SHOW 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 包含idendpoint(dnode 地址与 mqttPort 端口)、protocol(此处为mqtt)以及create_time等字段。

删除 Bnode

使用如下 SQL 语句删除 Bnode:

DROP 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.sqlsub.pyrun.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)、创建超级表、基于超级表创建 Topic(CREATE TOPIC语句将数据映射为可订阅的消息流)、插入一条测试数据,最后在 dnode 1 上创建 Bnode 启动订阅服务。

编写消费者

将以下 Python 代码保存为sub.py

import 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, properties=None): 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", qos=1, properties=sub_properties) def on_subscribe(client, userdata, mid, granted_qos, properties=None): 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_id="tmq_sub_cid", userdata=None, protocol=mqttClient.MQTTv5) else: client = mqttClient.Client(client_id="tmq_sub_cid", userdata=None, protocol=mqttClient.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-mqtt==2.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(虚拟节点组)ID
rows数据行数组,每行包含建 Topic 时 SELECT 出的全部列,如ts(毫秒时间戳)、tbname(子表名)、f1t1

进阶:集群化负载均衡与测试工具

共享订阅($share/group_id/topic)在多个消费者同属一个消费组时,消息会在组内分发,从而在多个消费客户端之间实现负载均衡与高可用;不同消费组之间则各自独立消费同一份数据。

如果需要压测或快速验证订阅链路,仓库 source/libs/tmqtt/tools 下提供了相关工具:

  • topic-producer.c:用于向 TDengine Topic 写入测试数据的生产者程序,其中同样配置了auto_offset_resetearliest/latest等消费参数;
  • perf.py:基于paho-mqtt的订阅压测脚本,演示了$share/g3/topic_meters等不同共享消费组的订阅写法。

小结

通过 Bnode 组件,TDengine 将内置的 Topic 数据流以标准 MQTT 协议对外暴露,使任意 MQTT 生态客户端都能直接消费时序数据,大幅降低了与消息中间件集成的门槛。实际部署时需注意:Topic 必须预先通过 SQL 创建;追求sub-offset等高级特性请使用 MQTT 5.0;mqttPort端口可在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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/12 17:18:58

渗透测试自学第十天:从HTTP协议到Burp Suite抓包改包实战

第十天&#xff0c;我没急着去学那些听起来很帅的东西&#xff0c;而是老老实实把HTTP协议的知识重新过了一遍&#xff0c;再把Burp Suite从安装到真正拦下第一个包&#xff0c;完整走通了一遍。这大概是自学渗透测试以来最踏实的一天——因为从这天开始&#xff0c;手头的工具…

作者头像 李华
网站建设 2026/9/12 17:18:55

模型蒸馏与数据版权争议:大模型 API 调用的合规边界与风险规避

1. 事件全景&#xff1a;一场关于“数据版权”的正面硬刚这两天 AI 圈炸了锅&#xff0c;Anthropic 直接在官网和社交媒体上公开点名了三大国产大模型&#xff0c;声称它们涉嫌“蒸馏”自家 Claude 系列模型的能力&#xff0c;而且证据相当具体。与此同时&#xff0c;马斯克也没…

作者头像 李华
网站建设 2026/9/12 17:18:33

多无人机协同路径规划的改进PSO算法与MATLAB实现

1. 项目背景与核心挑战多无人机协同作业已成为物流配送、农业植保、灾害救援等领域的重要技术手段。在复杂三维环境中实现多机动态避障路径规划&#xff0c;需要解决三个核心问题&#xff1a;实时环境感知与障碍物动态更新多目标优化&#xff08;路径长度、能耗、安全性等&…

作者头像 李华
网站建设 2026/9/12 17:12:21

DLA植物生长模拟:基于扩散凝聚的分形生成方法

简介&#xff1a;本资源是一套基于DLA&#xff08;扩散限制聚集&#xff09;算法的植物生长模拟程序&#xff0c;面向计算机图形学初学者、分形算法研究者及生物建模爱好者&#xff0c;用于理解分形几何与自然形态生成的内在关联。压缩包共22个文件&#xff0c;含5个C源码文件&…

作者头像 李华