ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

TDengine 数据订阅快速入门:用 Topic 与消费者组实现实时数据推送

TDengine 数据订阅快速入门:用 Topic 与消费者组实现实时数据推送 TDengine 数据订阅快速入门用 Topic 与消费者组实现实时数据推送【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine监控告警、实时分析、数据同步等场景中下游程序往往需要第一时间拿到新写入的数据。如果靠周期性轮询查询不仅延迟高还会给数据库带来额外负载。TDengine 内置了类似消息队列的数据订阅能力你可以用 SQL 把关注的数据定义成一个 TopicTDengine 会自动把持续写入的新数据按 Topic 推送给下游程序省去轮询逻辑也省去额外部署一套消息队列的复杂度。读完本文你将掌握如何创建查询 Topic、如何用taosshell 以消费者组订阅并实时查看推送结果、常用订阅参数的含义以及如何查看和清理订阅相关资源。本文以智能电表power库 /meters超级表模型为例与快速入门系列前几章保持一致关于订阅的完整语法与进阶主题可继续阅读 Data Subscription、Topic Syntax、Native Subscription、MQTT Data Subscription 与 Data Subscription API。订阅能力一览与 Kafka 类似在 TDengine 中你可以定义 Topic。一个 Topic 可以是一个数据库、一张超级表或者一段对现有表的查询过滤与预处理由 TDengine 在服务端完成。消费者可以加入消费者组共享消费进度。数据从 WAL 中推送采用**至少一次at-least-once**投递语义配合 ACK确认机制在崩溃、重启等故障场景下保证不丢消息。Topic 管理支持创建、查看、删除 Topic涵盖查询 Topic、超级表 Topic、数据库 Topic 三种类型并支持RELOAD TOPIC重新加载查询 Topic 定义。详见 Topic Syntax。消费者组与消费进度同一消费者组内的多个消费者共享消费进度不同消费者组即使订阅同一 Topic 也互不影响。可用SHOW CONSUMERS/SHOW SUBSCRIPTIONS查看状态用DROP CONSUMER GROUP清理。原生订阅通过各语言连接器 API 创建消费者、订阅 Topic、poll 并解析消息、提交 offset。详见 Native Subscription 和 Data Subscription API。MQTT 订阅自v3.3.7.0起MQTT 客户端可以连接 Bnodetaosmqtt服务直接订阅已有 Topic。详见 MQTT Data Subscription。下面我们从创建查询 Topic 开始用 shell 实时消费一遍完整流程。前置准备请先确认完成前面章节的两项前提TDengine 服务已启动并且能用taosshell 正常连接理解power数据库、meters超级表以及d1001、d1002等子表的基本模型。如果尚未创建这些对象在第一个 shell 中执行以下 SQL 完成初始化CREATE DATABASE IF NOT EXISTS power PRECISION ms KEEP 3650 DURATION 10 BUFFER 16; USE power; CREATE STABLE IF NOT EXISTS meters ( ts timestamp, current float, voltage int, phase float ) TAGS ( location varchar(64), group_id int ); CREATE TABLE IF NOT EXISTS d1001 USING meters TAGS (California.SanFrancisco, 2); CREATE TABLE IF NOT EXISTS d1002 USING meters TAGS (California.SanFrancisco, 3);创建 Topic在第一个 shell 中创建一个名为topic_meters的 Topic。Topic 决定了订阅者能收到哪些数据。下面这条 SQL 订阅meters超级表上新写入的数据并额外输出tbname这样你就能看到每一行来自哪张子表CREATE TOPIC IF NOT EXISTS topic_meters AS SELECT tbname, ts, current, voltage, phase FROM meters;执行以下命令确认 Topic 创建成功SHOW TOPICS;说明CREATE TOPIC ... AS subquery定义的是查询 Topic。它的数据粒度由 Topic SQL 决定——可以带过滤条件和标量函数但不支持聚合函数、时间窗口聚合以及DISTINCT、GROUP BY、ORDER BY、PARTITION BY、LIMIT/SLIMIT等子句订阅查询只能查询原始数据且只能按时间顺序输出。一个 TDengine 实例中 Topic 的最大数量由tmqMaxTopicNum控制默认 20详见 taosd 配置参数。打开第二个 Shell 并订阅另开一个终端进入taosshell执行订阅命令subscribe topic_meters -g quickstart_cg;其中topic_meters要订阅的 Topic 名称-g quickstart_cg指定消费者组。消费者组负责保存消费进度同一组再次订阅时会从已提交的位置继续消费。命令执行后 shell 会进入等待状态并显示类似如下的提示Subscribing to topic [topic_meters], group [quickstart_cg], offset [latest] ... Press CtrlC to stop.默认从latest最新位置开始订阅。保持这个 shell 不关闭回到第一个 shell 写入新数据。写入数据并查看订阅结果在第一个 shell 中插入两行新的电表数据INSERT INTO d1001 VALUES (NOW, 10.3, 219, 0.31); INSERT INTO d1002 VALUES (NOW, 10.2, 220, 0.23);在第二个 shell 中subscribe命令会实时输出新写入的数据。输出布局可能随终端宽度略有变化内容大致如下tbname | ts | current | voltage | phase | d1001 | 2026-07-24 18:20:01.000 | 10.3000 | 219 | 0.310 | d1002 | 2026-07-24 18:20:02.000 | 10.2000 | 220 | 0.230 |按CtrlC停止订阅。停止后 shell 会打印本次总共收到的行数Unsubscribed. Total rows received: 2从源码结构看这一交互由 shellSubscribe.c 实现它解析参数后调用 TMQ 客户端 API 依次完成tmq_conf_new创建配置、tmq_consumer_new创建消费者、tmq_subscribe订阅 Topic然后循环tmq_consumer_poll拉取消息按列宽打印表头与数据行收到-n指定的行数或按下CtrlC后调用tmq_unsubscribe与tmq_consumer_close清理资源见 shellSubscribe.c。也就是说shell 的subscribe命令本质上就是 TMQ 原生订阅接口的一个命令行最小实现非常适合开发、测试与排障时快速验证 Topic 是否能正常投递数据。常用订阅选项shell 订阅命令的完整格式为subscribe topic -g group_id [options];常用选项包括-c client_id指定客户端 ID默认自动生成-o earliest从最早可消费的位置开始适合消费 Topic 中已存在的存量数据-o latest从最新位置开始默认值适合实时等待新数据-n count收到指定行数后自动退出便于演示和测试-t timeout_mspoll 超时时间毫秒默认 1000。例如下面这条命令从最早位置开始读取收到 5 行后自动退出subscribe topic_meters -g quickstart_cg_earliest -o earliest -n 5;查看帮助subscribe -h;对应源码中 shellSubscribe.c 的帮助文本还给出了更多用法示例例如subscribe my_topic -g group1 -t 500;用于自定义 poll 超时。查看与清理订阅资源在 shell 中可以查看 Topic、消费者与订阅分配情况SHOW TOPICS; SHOW CONSUMERS; SHOW SUBSCRIPTIONS;SHOW TOPICS显示当前数据库下所有 Topic 的信息完整字段见元数据表INS_TOPICSSHOW CONSUMERS显示当前数据库下所有消费者信息含状态与创建时间完整字段见性能表PERF_CONSUMERSSHOW SUBSCRIPTIONS显示一个 Topic 在各 vgroup 上的消费情况便于监控消费进度完整字段见元数据表INS_SUBSCRIPTIONS。当不再需要本快速入门示例时先停止订阅 shell再执行以下 SQL 清理DROP CONSUMER GROUP IF EXISTS FORCE quickstart_cg ON topic_meters; DROP CONSUMER GROUP IF EXISTS FORCE quickstart_cg_earliest ON topic_meters; DROP TOPIC IF EXISTS topic_meters;关于清理的几点说明单个消费者无法单独删除但可以删除其所属的消费者组DROP CONSUMER GROUP [IF EXISTS] [FORCE] cgroup_name ON topic_name;。如果组内还有活跃消费者需要加FORCE强制删除强制删除后这些消费者在消费时会报错FORCE自v3.3.6.0起支持。删除 Topic 同理DROP TOPIC [IF EXISTS] [FORCE] topic_name;若仍有消费者订阅需用FORCE强制删除。深入三种 Topic 与更多消费方式快速入门演示的是最常用的查询 Topic。除此之外TDengine 还支持另外两种通过 SQL 创建的 Topic详见 Topic Syntax超级表 TopicCREATE 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只能过滤 tag 或tbname而不能使用普通列。数据库 TopicCREATE TOPIC [IF NOT EXISTS] topic_name [WITH META | ONLY META] AS DATABASE db_name;订阅数据库内所有表的数据。超级表与数据库订阅属于高级模式使用时请谨慎并参考技术资料。消费方式也有三种Data Subscriptionshell 订阅即本文演示的方式适合快速验证原生订阅通过 C、Java、Go、Rust、Python、C# 等语言连接器 API 创建消费者td.connect.ip、group.id、auto.offset.reset、enable.auto.commit、auto.commit.interval.ms、session.timeout.ms、max.poll.interval.ms、fetch.max.wait.ms、min.poll.rows等参数调用 subscribe/poll/commit 完成消费闭环Native Subscription自v3.2.0.0起订阅还支持 vnode 迁移与拆分迁移/拆分前请先把 WAL 中的存量数据消费完MQTT 订阅自v3.3.7.0起创建 BnodeCREATE BNODE ON DNODE dnode_id默认端口 6057可用mqttPort参数修改后MQTT 客户端即可直接订阅已有 Topic支持$share/group_id/topic_name共享订阅、sub-offsetearliest指定起始位置以及 QoS 0/1MQTT Data Subscription。此外TDengine 订阅还支持数据回放Replay通过消费者参数enable.replaytrue开启后消息会按原始写入时的时间间隔重新推送便于以原始节奏重跑数据流只有查询 Topic 支持回放回放进度不会被保存。所有订阅均基于 WALTDengine 自动为 WAL 文件建立随机访问索引并提供可配置的轮转与保留策略保留时间与大小使 WAL 成为持久化、保持到达顺序的存储引擎查询 Topic 从 WAL 读取后由统一查询引擎按当前 offset 做过滤与变换再推送给消费者。进一步阅读本章只覆盖了用 shell 快速验证查询 Topic 订阅的常用流程。完整的订阅能力请继续阅读Data Subscription数据订阅总览Topic 与消费者组、WAL、消费模型Topic SyntaxCREATE/DROP/SHOW TOPIC、三种 Topic 类型、消费者组与回放注意事项Native Subscription通过连接器 API 创建消费者并订阅MQTT Data Subscription用 MQTT 客户端连接 Bnode 订阅 Topic 数据Data Subscription API多语言连接器订阅 API 与完整示例【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表