完全指南:主题、消费组、Native 与 MQTT 订阅实战)
数据库时序数据库物联网大数据实时分析云原生【免费下载链接】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 在时序数据库内核中内置了与 Kafka 高度相似的消息队列能力用户可以用 SQL 定义主题topic下游应用通过多语言连接器或 MQTT 客户端实时消费新写入的数据从而在监控、告警、实时分析与数据同步等场景中省去额外部署消息队列的成本。本文基于仓库中的数据订阅章节及其配套文档系统讲解主题语法、消费组与 WAL 存储原理、Native 连接器消费流程以及 MQTT 订阅方案并给出可直接运行的 SQL、Python 示例与源码级佐证。为什么 TDengine 需要内置数据订阅在监控、告警、实时分析和数据同步等场景中下游程序通常需要第一时间获取新写入的数据。如果通过定时查询拉取数据不仅延迟更高也会增加数据库查询压力。TDengine 提供了类似于消息队列产品的数据订阅和消费接口。在许多场景中采用 TDengine 的时序大数据平台无须再集成消息队列产品从而简化应用程序设计并降低运维成本。与 Kafka 类似用户需要在 TDengine 中定义主题topic。TDengine 的主题可以是一个数据库、一张超级表或者基于现有超级表、子表或普通表的查询条件即一条查询语句。用户可以利用 SQL 对标签、表名、列、表达式等条件进行过滤并对数据进行标量函数与 UDF 计算不包括数据聚合。与其他消息队列工具相比这是 TDengine 数据订阅功能的最大优势数据的粒度由定义主题的 SQL 决定过滤与预处理由 TDengine 自动完成从而减少传输的数据量并降低应用程序的复杂度。核心概念主题、消费组与 WAL消费者与消费组消费者订阅主题后可以实时接收最新的数据。多个消费者可以组成一个消费组共享消费进度实现多线程、分布式地消费数据提高消费速度。不同消费组的消费者即使消费同一个主题也不共享消费进度。一个消费者可以订阅多个主题。如果主题对应的是超级表或库数据可能会分布在多个不同的节点或数据分片上当一个消费组中有多个消费者时可以提高消费效率。TDengine 的消息队列提供了消息的 ACKAcknowledgment确认机制确保在宕机、重启等复杂环境下实现至少一次at least once消费。WAL把预写日志改造为持久化消息存储为实现上述功能TDengine 会为预写数据日志Write-Ahead LoggingWAL文件自动创建索引以支持快速随机访问并提供了灵活可配置的文件切换与保留机制。用户可以根据需求指定 WAL 文件的保留时间和大小。通过这些方法WAL 被改造成一个保留事件到达顺序的、可持久化的存储引擎。对于以主题形式创建的查询TDengine 将从 WAL 读取数据。在消费过程中TDengine 根据当前消费进度从 WAL 直接读取数据并使用统一的查询引擎实现过滤、变换等操作然后将数据推送给消费者。从源码结构看WAL 的保留策略在 vnode 层被显式配置与校验。例如 vnodeSvr.c 中处理 vnode 更新请求时会对比walRetentionPeriod、walRetentionSize等字段并写入 vnode 配置相关配置参数的说明可参考 taosd 配置参数。主题数量的实例级上限一个 TDengine 实例可创建的 topic 个数上限由tmqMaxTopicNum参数控制。该参数定义于 tglobal.c默认值 20取值范围 110000属全局配置参数支持通过 SQL 动态修改、立即生效并在 mnode 侧创建 topic 时被实际校验——mndTopic.c 中通过sdbGetSize(pMnode-pSdb, SDB_TOPIC) tmqMaxTopicNum判断是否达到上限。完整参数说明见 taosd 配置参数 · tmqMaxTopicNum。版本相关行为变更从v3.2.0.0开始数据订阅支持 vnode 迁移和分裂。由于数据订阅依赖 WAL 文件而在 vnode 迁移和分裂过程中WAL 文件并不会同步。因此在迁移或分裂操作完成后将无法继续消费此前尚未消费完的 WAL 数据。请务必在执行 vnode 迁移或分裂之前将所有 WAL 数据消费完毕。从v3.3.7.0开始支持通过 MQTT 客户端订阅已创建主题的数据详见 MQTT 订阅。主题类型与创建语法TDengine 从v3.0.0.0开始对消息队列做了大幅优化和增强以简化用户的数据订阅方案。用户可以通过 SQL 创建订阅主题然后使用连接器 API、taosshell 或 MQTT 客户端消费主题中的数据。TDengine 使用 SQL 创建的主题共有 3 种类型。查询主题SELECT 主题订阅一条 SQL 查询定义的数据流创建语法如下CREATE TOPIC [IF NOT EXISTS] topic_name AS subquery该 SQL 通过SELECT语句订阅数据包括SELECT *或SELECT ts, c1等指定列查询。查询主题可以带条件过滤和标量函数计算但不支持聚合函数、时间窗口聚合也不支持DISTINCT、GROUP BY、ORDER BY、PARTITION BY、LIMIT/SLIMIT等。需要注意该类型 topic 一旦创建订阅数据的结构即确定。被订阅或用于计算的列或标签不可被删除ALTER TABLE DROP或修改ALTER TABLE MODIFY。从v3.4.0.0开始可以修改、删除、增加这些列或标签但需要执行RELOAD TOPIC使变更生效。对于SELECT *订阅会展开为创建 topic 时的所有列子表、普通表为数据列超级表为数据列加标签列。不支持虚拟表的查询订阅。subquery中的超级表、子表、普通表可以被删除。删除后订阅数据为空如果删除后重新创建同名表订阅数据仍然为空因为表 ID 已变化。如需订阅新建表的数据可以通过RELOAD TOPIC重新加载 topic。实战示例假设需要订阅所有智能电表中电压值大于 200 的数据且仅仅返回时间戳、电流、电压 3 个采集量不返回相位可以通过下面的 SQL 创建power_topic主题CREATE TOPIC power_topic AS SELECT ts, current, voltage FROM power.meters WHERE voltage 200;超级表主题订阅一个超级表中的所有数据语法如下CREATE TOPIC [IF NOT EXISTS] topic_name [WITH META | ONLY META] AS STABLE stb_name [where_condition]与使用SELECT * FROM stbName订阅的区别是不会限制用户的表结构变更即表结构变更以及变更后的新数据都能够订阅到。返回的是非结构化的数据返回数据的结构会随着超级表的表结构变化而变化。WITH META参数可选指定后将返回创建超级表、子表等语句主要用于 taosX 做超级表迁移。ONLY META参数可选指定后仅订阅元数据变更不再传输时序数据。where_condition参数可选用于过滤符合条件的子表并订阅这些子表。WHERE条件里不能有普通列只能是 tag 或tbname可以使用函数过滤 tag但不能使用聚合函数因为子表 tag 值无法做聚合也可以是常量表达式比如2 1订阅全部子表或者false订阅 0 个子表。返回数据不包含标签。支持虚拟超级表的订阅仅能订阅出虚拟超级表的 meta 信息所以虚拟超级表订阅需要带上WITH META或ONLY META参数否则订阅不到内容。数据库主题订阅一个数据库里所有数据语法如下CREATE TOPIC [IF NOT EXISTS] topic_name [WITH META | ONLY META] AS DATABASE db_name;通过该语句可创建一个包含数据库所有表数据的订阅WITH META参数可选指定后将返回数据库里所有超级表、子表、普通表的元数据创建、删除、修改语句主要用于 taosX 做数据库迁移。ONLY META参数可选指定后仅订阅元数据变更不再传输时序数据。WITH META或ONLY META的情况下可订阅出虚拟表的信息并且仅能订阅出虚拟表的 meta 信息。说明超级表订阅和库订阅属于高级订阅模式容易出错如确实要使用请咨询技术支持人员。主题管理删除、查看与加载删除主题如果不再需要订阅数据可以删除 topic。如果当前 topic 被消费者订阅通过FORCE语法可强制删除强制删除后订阅的消费者在消费数据时会出错FORCE语法从v3.3.6.0开始支持。DROP TOPIC [IF EXISTS] [FORCE] topic_name;查看主题SHOW TOPICS;显示当前数据库下的所有主题信息。更完整字段见元数据表INS_TOPICS。加载主题RELOAD TOPICRELOAD TOPIC [IF EXISTS] topic_name AS subquery;该语法从v3.4.0.0开始支持仅适用于查询主题用于重新加载主题定义。它主要解决查询主题里变更列或 tag以及SELECT *查询订阅时删除或增加列、tag 后输出结果不生效的问题。需要变更订阅表结构的 schema 时建议先停止消费再变更表结构然后执行RELOAD TOPIC接着重新开始订阅。消费者与消费组管理创建消费者消费者通常通过 TDengine 客户端驱动或者连接器所提供的 API 创建详情可以参考开发指南 · 数据订阅。为了快速验证订阅功能也可以在taosshell 中执行subscribe topic -g group_id创建消费者并开始消费具体用法请参考taosCLI 数据订阅。查看消费者SHOW CONSUMERS;显示当前数据库下所有消费者的信息包括消费者的状态、创建时间等。更完整字段见性能数据表PERF_CONSUMERS。删除消费组创建消费者时会为其指定一个消费者组。消费者不能显式地删除但可以删除消费者组。如果当前消费者组里有消费者在消费通过FORCE语法可强制删除强制删除后订阅的消费者在消费数据时会出错FORCE语法从v3.3.6.0开始支持。DROP CONSUMER GROUP [IF EXISTS] [FORCE] cgroup_name ON topic_name;查看订阅信息SHOW SUBSCRIPTIONS;显示 topic 在不同 vgroup 上的消费信息可用于查看消费进度。更完整字段见元数据表INS_SUBSCRIPTIONS。Native 订阅用连接器 API 消费数据TDengine 提供了多语言数据订阅 API包括但不限于创建消费者、订阅主题、取消订阅、拉取数据、提交与设置消费进度等与 Kafka 订阅 API 保持高度一致便于复用既有开发经验。目前支持 C、Java、Go、Rust、Python 和 C# 等。各语言连接器对SELECT主题、库/超级表主题以及WITH META/ONLY META消费的支持差异详见数据订阅编程接口中的能力对照表。创建主题如下 SQL 将创建一个名为topic_meters的订阅。使用该订阅所获取的消息中的每条记录都由查询语句SELECT ts, current, voltage, phase, groupid, location FROM meters所选择的列组成。CREATE TOPIC IF NOT EXISTS topic_meters AS SELECT ts, current, voltage, phase, groupid, location FROM meters;创建消费者的常用参数TDengine 消费者的概念与 Kafka 类似消费者通过订阅主题来接收数据流。消费者可以配置多种参数如连接方式、服务器地址、自动提交 Offset、自动重连、数据传输压缩等。常用参数如下参数说明默认值td.connect.ip服务端的 FQDN—td.connect.user用户名—td.connect.pass密码—td.connect.tokenToken通过CREATE TOKEN语句生成使用 Token 认证时无需设置 user/pass—td.connect.port服务端的端口号—group.id消费组 ID同一消费组共享消费进度必填项最大长度 192不可包含英文冒号:每个 topic 最多可建立 100 个 consumer group—client.id客户端 ID最大长度 255—auto.offset.reset消费组订阅的初始位置earliest从头订阅v3.2.0.0之前默认latest仅从最新数据订阅v3.2.0.0及之后默认none没有已提交 offset 时无法订阅latestenable.auto.commit是否启用消费位点自动提交false 时需应用自行 committrueauto.commit.interval.ms自动提交消费位点的时间间隔毫秒5000msg.with.table.name是否允许从消息中解析表名不适用于列订阅自v3.2.0.0起废弃关闭enable.replay是否开启数据回放功能关闭session.timeout.ms消费者心跳丢失后的超时时间v3.3.3.0起支持范围[6000, 1800000]超时触发 rebalance 后该 consumer 会被删除12000max.poll.interval.ms消费者拉取数据间隔的最长时间v3.3.3.0起支持范围[1000, INT32_MAX]超过则视为离线并触发 rebalance300000fetch.max.wait.ms服务端单次返回数据的最大耗时v3.3.6.0起支持范围[1, INT32_MAX]1000min.poll.rows服务端单次返回数据的最小条数v3.3.6.0起支持范围[1, INT32_MAX]4096高级参数tmq_conf_new默认均为关闭enable.wal.marker提交位点时是否向 mnode 发送 WAL markerboolean默认false。仅当提交的 offset 类型为 WAL log 时生效一般用于迁移类工具。msg.enable.batchmeta是否启用服务端批量元数据返回多条 meta 可合并为 batch meta 响应非0开启默认关闭。Java WebSocket 侧属性名为enable_batch_meta。msg.consume.rawdata消费数据时拉取数据类型为二进制类型不可做解析操作内部参数只用于 taosX 数据迁移默认0不起效非0起效v3.3.6.0起支持。各语言连接器还提供各自专属参数如 Go 的ws.url支持多节点故障切换如ws://node1:6041,ws://node2:6041、timezoneIANA 时区Go 连接器3.7.4及以上支持等完整列表见数据订阅编程接口。订阅与消费流程消费者订阅主题后可以开始接收并处理这些主题中的消息。典型流程如下订阅数据调用订阅接口指定主题列表名称支持同时订阅多个主题。拉取数据调用 poll 类接口每次调用获取一条消息一条消息中可能包含多条记录。解析结果按各语言连接器约定解析消息字段字段名和数据类型与主题定义的列一一对应。限制提醒订阅查询只能查询原始数据不能查询聚合或计算结果订阅查询只能按照时间正序查询数据。指定订阅的 Offset 与提交 Offset消费者可以指定从特定 Offset 开始读取分区中的消息从而重读消息或跳过已处理的消息。当消费者读取并处理完消息后可以提交 Offset表示已经成功处理到该 Offset。Offset 提交可以是自动的根据配置定期提交或手动的由应用程序控制何时提交。当创建消费者时属性enable.auto.commit为 false 时可以手动提交 offset如 C 语言的tmq_commit_sync、Rust 的consumer.commit。手工提交前确保消息正常处理完成否则处理出错的消息不会被再次消费自动提交是在本次poll消息时可能提交上次消息的消费进度因此请确保消息处理完毕再进行下一次poll。取消订阅和关闭消费消费者可以取消对主题的订阅停止接收消息。当消费者不再需要时应关闭消费者实例以释放资源并断开与 TDengine 服务器的连接。注意消费者取消订阅后已经关闭无法重用如果想订阅新的 topic请重新创建消费者。各语言的完整可运行示例WebSocket 与原生连接两种方式见开发指南 · 数据订阅及仓库中的示例代码如 c/tmq_demo.c、python/tmq_native.py、go/tmq/native/main.go 等。回放功能ReplayTDengine 的数据订阅支持回放按数据实际写入时间间隔重新推送消息便于按原有节奏重放数据流。该能力基于 WAL 实现。例如写入如下 3 条数据时回放会先返回第 1 条约 5s 后返回第 2 条再约 3s 后返回第 3 条2023/09/22 00:00:00.000 2023/09/22 00:00:05.000 2023/09/22 00:00:08.000使用回放功能时需要注意通过消费参数enable.replay为true开启回放。仅查询主题支持回放超级表主题和数据库主题不支持。回放不支持进度保存。回放本身需要处理时间精度存在约数十毫秒的误差。主题 SQL 中可用WHERE限定时间范围或过滤条件这属于主题定义本身与是否开启回放无关。MQTT 订阅通过 Bnode 使用标准 MQTT 客户端TDengine 从v3.3.7.0开始提供 MQTT 订阅功能。通过 MQTT 客户端连接 TDengine Bnode 服务可直接订阅系统中已有主题的数据。主要特性协议支持推荐使用 MQTT 5.0亦兼容 MQTT 3.1 / 3.1.1sub-offset等用户属性依赖 MQTT 5.0。身份验证使用 TDengine 原生验证。主题管理与标准 MQTT 协议不同主题必须预先创建因不支持消息发布无法通过发布消息动态创建。共享主题形如$share/group_id/topic_name的主题被视为共享订阅适用于需要负载均衡和高可用的场景。订阅位置支持latest、earliestWAL 最早位置可通过订阅用户属性sub-offsetearliest指定默认latest。服务质量支持 QoS 0、QoS 1。Bnode 节点管理用户可通过 TDengine 的命令行工具taos管理 Bnode。执行命令前请确保taos可正常连接集群。创建 BnodeCREATE BNODE ON DNODE {dnode_id}一个 Dnode 上只能创建一个 Bnode。Bnode 创建成功后会自动启动 Bnode 子进程taosmqtt默认在 6057 端口对外提供 MQTT 订阅服务端口可在taos.cfg中通过参数mqttPort配置默认值 6057取值范围 165056局部配置参数不支持动态修改详见 taosd 配置参数 · mqttPort。该参数在 tglobal.c 中被注册为mqttPort配置项。例如CREATE BNODE ON DNODE 1。查看 Bnode列出集群中所有的数据订阅节点包括其id、endpoint、create_time等属性。更完整字段见元数据表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)删除 BnodeDROP BNODE ON DNODE {dnode_id}删除 Bnode 将把 Bnode 从 TDengine 集群中移除同时停止taosmqtt服务。订阅数据示例环境准备在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;客户端订阅可以使用兼容 MQTT 协议的客户端来订阅前一步环境中的数据这里使用 Pythonpaho-mqtt举例说明示例按 MQTT 5.0 编写以便设置sub-offset。在操作系统命令行中依次执行下面这些命令即可订阅到上一步写入的数据订阅成功后若topic_meters主题中有新增写入则会通过 MQTT 协议推送到客户端。python3 -m venv .test-env source .test-env/bin/activate pip3 install paho-mqtt2.1.0 python3 ./sub.py其中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, 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()消息格式上一节的示例中会输出下面的信息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 值后面是 UTF-8 编码的 JSON 消息其中rows是数据行的数组。深入原理与进一步阅读从源码实现看主题topic作为 mnode 管理的一等元数据对象被持久化在 SDB 中mndTopic.c中定义了SDB_TOPIC存储类型及 topic 的序列化/反序列化逻辑创建 topic 时即通过sdbGetSize校验实例级数量上限mndTopic.c。订阅数据的读取则发生在 vnode 侧消费时依据消费进度从 WAL 直接读取再交给统一查询引擎做过滤与变换这与文档中WAL 被改造成保留事件到达顺序的可持久化存储引擎的描述相互印证。如需继续深入可以依次阅读主题语法三种主题类型、WITH META/ONLY META语义与RELOAD TOPIC细节。Native 订阅消费者参数与消费流程的概念说明。MQTT 订阅Bnode 管理、MQTT 消息格式与完整示例。数据订阅编程接口各语言连接器C、Java、Go、Rust、Python、C#、Node.js的创建参数、WebSocket 与原生连接的完整代码示例。元数据与性能表INS_TOPICS、INS_SUBSCRIPTIONS、INS_BNODES、PERF_CONSUMERS可用于查看主题、消费进度与消费者状态。赞分享数据库时序数据库物联网大数据实时分析云原生【免费下载链接】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 配置到 paho-mqtt 消费端实战TDengine MQTT 数据订阅完全指南从 Bnode 配置到 paho mqtt 消费端实战 从 v3.3.7.0 版本起TDengine 原生支持通数据库时序数据库物联网大数据实时分析云原生TDengine 数据订阅引擎内幕TMQ Topic、消费者组与 Rebalance 机制TDengine 数据订阅引擎内幕TMQ Topic、消费者组与 Rebalance 机制 TDengine 的数据订阅TMQTDengine Messa数据库时序数据库大数据物联网云原生TDengine数据订阅(TMQ)功能详解TDengine数据订阅 TMQ 功能详解 概述 TDengine的TMQ 时序消息队列 功能提供了一种高效的数据订阅机制允许应用程序实时获取数据库中的数据变数据库时序数据库大数据物联网云原生上一篇TradingAgents-CN如何用多智能体AI框架构建你的智能投资分析系统下一篇终极指南解决Azure AKS中Node自动供应(NAP)功能因ipFamilies参数导致的启用失败问题创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考