
数据库时序数据库物联网大数据实时分析云原生【免费下载链接】tdengineTDengine 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 从v3.3.7.0起内置 MQTT 订阅能力只需启动一个 BnodeBroker 节点任何兼容 MQTT 的客户端即可像连接普通 Broker 一样订阅 TDengine 中预先创建的主题实时接收数据变更。本文以官方文档 MQTT 订阅 为主体结合仓库中 Bnode 与 tmqtt 模块的源码实现完整讲解 Bnode 的创建、查看、删除MQTT 客户端的连接与订阅以及消息格式与订阅位置offset等核心机制帮助读者从零搭建一条SQL 写入 → 主题推送 → 任意 MQTT 客户端消费的数据订阅链路。功能概述与核心特性TDengine 的 MQTT 订阅功能定位非常明确通过 MQTT 客户端连接 TDengine Bnode 服务直接订阅系统中已有主题的数据。它与标准 MQTT Broker 的关键差异在于TDengine 不承担消息发布职责主题是数据库内真实存在的数据集合因此主题必须预先通过 SQL 创建无法像普通 Broker 那样通过发布消息动态建主题。该功能具备以下核心特性特性说明协议支持推荐使用 MQTT 5.0同时兼容 MQTT 3.1 / 3.1.1其中sub-offset等用户属性依赖 MQTT 5.0身份验证使用 TDengine 原生账号体系验证示例默认账号密码为root/taosdata主题管理主题必须预先创建CREATE TOPIC不支持通过发布消息动态创建共享主题形如$share/group_id/topic_name的主题被视为共享订阅适用于负载均衡和高可用场景订阅位置支持latest默认与earliestWAL 最早位置通过订阅用户属性sub-offsetearliest指定服务质量支持 QoS 0至多一次与 QoS 1至少一次从源码结构看Bnode 是 MQTT 能力的承载单元source/dnode/bnode/src/bnode.c 中的bndOpen()在检测到节点协议为TSDB_BNODE_OPT_PROTO_MQTT时会调用mqttMgmtStartMqttd()拉起底层 MQTT 服务进程因此Bnode 即 MQTT Broker是本功能最直观的架构认知。Bnode 节点管理用户可通过 TDengine 的命令行工具taos管理 Bnode。执行下述命令前请确保taos可正常连接集群。创建 BnodeCREATE BNODE ON DNODE {dnode_id}关键约束与行为一个 Dnode 上只能创建一个 BnodeBnode 创建成功后会自动启动 Bnode 子进程taosmqtt默认在6057端口对外提供 MQTT 订阅服务端口可在taos.cfg中通过参数mqttPort配置详见 taosd 配置参数。例如在 dnode 1 上创建 BnodeCREATE BNODE ON DNODE 1;关于mqttPort的取值source/common/src/tglobal.c 中的配置注册代码显示其为整型参数取值范围为1 ~ 65056cfgAddInt32(pCfg, mqttPort, defaultMqttPort, 1, 65056, ...)超出该范围的值将无法通过配置校验这点在自定义端口时需特别注意。查看 BnodeSHOW BNODES用于列出集群中所有的数据订阅节点包括其id、endpoint、create_time等属性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)其中endpoint即为 MQTT 客户端需要连接的地址与端口protocol字段标识该节点的协议类型为mqtt。若需查询更完整的字段如 dnode 关联信息、状态等可查看元数据表INS_BNODES。删除 BnodeDROP BNODE ON DNODE {dnode_id}删除 Bnode 会将 Bnode 从 TDengine 集群中移除同时停止taosmqtt服务此后该节点将不再对外提供 MQTT 订阅能力。订阅数据示例下面通过一个端到端示例演示如何用 Pythonpaho-mqtt客户端订阅 TDengine 主题并实时接收数据。该示例对应仓库中的 source/libs/tmqtt/example/sub.py可直接参考运行。环境准备首先在taos命令行中依次执行下面的 SQL创建数据库、超级表、主题topic_meters、Bnode并写入一条数据供下一步订阅使用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;这里值得说明的是CREATE TOPIC创建的topic_meters本质上是绑定在超级表db.meters上的一个数据投影包含ts、tbname、f1、t1四列。后续任何写入到db.meters的新数据都会实时反映到该主题中这正是主题必须预先创建这一特性的数据基础。客户端订阅在操作系统命令行中依次执行下面这些命令即可订阅到上一步写入的数据订阅成功后若topic_meters主题中有新增写入则会通过 MQTT 协议推送到客户端python3 -m venv .test-env source .test-env/bin/activate pip3 install paho-mqtt2.1.0 python3 ./sub.py其中sub.py文件的内容如下示例按 MQTT 5.0 编写以便设置sub-offset用户属性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, 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()示例关键点拆解连接地址与端口client.connect(127.0.1.1, 6057)中的端口 6057 即 Bnode 的 MQTT 服务端口与SHOW BNODES输出中的endpoint保持一致若通过mqttPort修改了端口此处需同步修改。身份认证client.username_pw_set(root, taosdata)使用 TDengine 原生账号密码完成 MQTT 连接验证与登录taos的凭据一致。共享订阅订阅主题写为$share/g1/topic_meters其中g1是分组名。多个客户端若使用相同的group_id此处为g1消息会在这些客户端之间负载均衡分发从而实现水平扩展与高可用。订阅位置控制通过 MQTT 5.0 的SUBSCRIBE包用户属性sub-offsetearliest指定从 WAL 最早位置开始消费不设置该属性时默认latest只消费订阅之后新写入的数据。由于该属性依赖 MQTT 5.0 用户属性机制使用 MQTT 3.1/3.1.1 的客户端将无法指定订阅位置。QoS示例使用qos1保证消息至少送达一次如对可靠性要求不高可改为qos0以获取更低的延迟与开销。订阅位置的底层实现从源码看sub-offset最终映射为底层消费者配置的auto_offset_reset。source/libs/tmqtt/mqtt/src/tmqttCtx.c 中tmq_ctx_topic_exists()构造消费者配置时auto_offset_reset earliest ? earliest : latest随后基于 TDengine 自研 tmq 消费组件完成真正的数据拉取。也就是说MQTT 订阅并非简单转发而是将 MQTT 协议适配到 TDengine 原生的主题消费通道之上这保证了订阅位置的精确性与 WAL 级别的数据一致性。消息格式上一节的示例运行后会输出下面的信息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本次消息所属的订阅主题1该条消息的 QoS 值对应示例中qos1其余部分为UTF-8 编码的 JSON 消息体各字段含义如下字段含义topic消息来源主题名db主题所在数据库vid数据所属的虚拟节点vnodeIDrows数据行数组每行包含该主题投影的列示例中为ts、tbname、f1、t1其中ts为毫秒级时间戳tbname为子表名。由于消息体是标准 JSON因此任意语言Python、Java、Go、Node.js 等的 MQTT 客户端都可以轻松完成反序列化与业务对接无需引入 TDengine 专有 SDK。典型应用场景与使用建议综合文档与实现MQTT 订阅适合以下场景异构系统集成将 TDengine 中的数据变更以标准 MQTT 协议推送给 Kafka、消息中间件、边缘网关或第三方应用避免为每种消费者引入专有驱动负载均衡消费利用$share/group_id/topic_name共享订阅让一组消费者实例分摊消息处理压力配合DROP BNODE/CREATE BNODE完成节点级的高可用切换回放与追平需要从历史最早位置重新消费时使用 MQTT 5.0 用户属性sub-offsetearliest依赖 WAL 中保留的数据具体保留时长受 WAL 配置影响。实操层面的建议主题与表结构强绑定规划订阅前应先用CREATE TOPIC ... AS SELECT ...精确限定需要投影的列多消费者组共享同一主题时注意区分group_id避免相互抢占消息生产环境务必使用 TDengine 原生账号体系配置独立账号与最小权限而不是沿用示例中的默认root账号。小结TDengine 的 MQTT 订阅功能以 Bnode 为桥接节点将数据库内预先创建的主题通过标准 MQTT 5.0 协议对外开放兼顾了协议通用性兼容 3.1/3.1.1与 TDengine 原生能力原生认证、WAL 级订阅位置、tmq 消费通道。本文覆盖了 Bnode 生命周期管理CREATE/SHOW/DROP BNODE、mqttPort端口配置、paho-mqtt 客户端订阅的完整流程以及 JSON 消息格式读者可基于此快速构建跨系统的时间序列数据订阅链路。赞分享数据库时序数据库物联网大数据实时分析云原生【免费下载链接】tdengineTDengine 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 MQTT 订阅指南通过 Bnode 与 taosmqtt 将主题数据推送给 MQTT 客户端TDengine MQTT 订阅指南通过 Bnode 与 taosmqtt 将主题数据推送给 MQTT 客户端 TDengine 从 v3.3.7.0 开始提数据库时序数据库大数据物联网云原生TDengine MQTT 数据订阅完全指南从 Bnode 配置到 paho-mqtt 消费端实战TDengine MQTT 数据订阅完全指南从 Bnode 配置到 paho mqtt 消费端实战 从 v3.3.7.0 版本起TDengine 原生支持通数据库时序数据库物联网大数据实时分析云原生TDengine 基于 MQTT 的数据订阅Bnode 管理与 taosmqtt 消费实践TDengine 基于 MQTT 的数据订阅Bnode 管理与 taosmqtt 消费实践 自 v3.3.7.0 起TDengine 原生支持通过 MQTT数据库时序数据库大数据物联网云原生创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考