数据订阅
在监控、告警、实时分析和数据同步等场景中,下游程序通常需要第一时间获取新写入的数据。如果通过定时查询拉取数据,不仅延迟更高,也会增加数据库查询压力。TDengine 提供了类似于消息队列产品的数据订阅和消费接口。在许多场景中,采用 TDengine 的时序大数据平台,无须再集成消息队列产品,从而简化应用程序设计并降低运维成本。
与 Kafka 类似,用户需要在 TDengine 中定义主题(topic)。TDengine 的主题可以是一个数据库、一张超级表,或者基于现有超级表、子表或普通表的查询条件,即一条查询语句。用户可以利用 SQL 对标签、表名、列、表达式等条件进行过滤,并对数据进行标量函数与 UDF 计算(不包括数据聚合)。与其他消息队列工具相比,这是 TDengine 数据订阅功能的最大优势:数据的粒度由定义主题的 SQL 决定,过滤与预处理由 TDengine 自动完成,从而减少传输的数据量并降低应用程序的复杂度。
消费者订阅主题后,可以实时接收最新的数据。多个消费者可以组成一个消费组,共享消费进度,实现多线程、分布式地消费数据,提高消费速度。不同消费组的消费者即使消费同一个主题,也不共享消费进度。一个消费者可以订阅多个主题。如果主题对应的是超级表或库,数据可能会分布在多个不同的节点或数据分片上;当一个消费组中有多个消费者时,可以提高消费效率。TDengine 的消息队列提供了消息的 ACK(Acknowledgment,确认)机制,确保在宕机、重启等复杂环境下实现至少一次(at least once)消费。
为实现上述功能,TDengine 会为预写数据日志(Write-Ahead Logging,WAL)文件自动创建索引,以支持快速随机访问,并提供了灵活可配置的文件切换与保留机制。用户可以根据需求指定 WAL 文件的保留时间和大小。通过这些方法,WAL 被改造成一个保留事件到达顺序的、可持久化的存储引擎。对于以主题形式创建的查询,TDengine 将从 WAL 读取数据。在消费过程中,TDengine 根据当前消费进度从 WAL 直接读取数据,并使用统一的查询引擎实现过滤、变换等操作,然后将数据推送给消费者。
从 v3.2.0.0 开始,数据订阅支持 vnode 迁移和分裂。由于数据订阅依赖 WAL 文件,而在 vnode 迁移和分裂过程中,WAL 文件并不会同步。因此,在迁移或分裂操作完成后,将无法继续消费此前尚未消费完的 WAL 数据。请务必在执行 vnode 迁移或分裂之前,将所有 WAL 数据消费完毕。
从 v3.3.7.0 开始,还支持通过 MQTT 客户端订阅已创建主题的数据,详见 MQTT 订阅。








