数据订阅
数据订阅
1. 功能介绍
IoTDB 数据订阅模块(下称 IoTDB 订阅客户端)为用户提供了一种区别于数据查询的流式数据消费方式。它参考了 Kafka 等消息队列产品的基本概念和逻辑,提供数据订阅和消费接口,但并不是为了完全替代这些消费队列的产品,更多的是在简单流式获取数据的场景为用户提供更加便捷的数据订阅服务。
在下面应用场景中,使用 IoTDB 订阅客户端消费数据会有显著的优势:
- 实时获取新增数据:订阅客户端可以持续拉取新写入的数据,无需频繁执行定时查询,从而降低应用开发复杂度和系统查询负载,适用于大屏展示、组态监控等需要及时刷新数据的场景。
- 便捷对接第三方系统:Flink、Kafka、DataX、MySQL 等下游系统可以作为订阅客户端主动拉取数据,无需在 IoTDB 内部为不同系统分别开发数据推送组件,能够简化系统集成和数据流转链路。
- 增量备份 TsFile:通过订阅新生成的 TsFile,用户可以及时将增量数据归档到指定存储位置,适用于定期备份和异地保存等场景。
- 低延迟实时计算:通过共识订阅直接消费已经进入 IoTConsensus 共识写入链路的数据,减少传统订阅链路对 Pipe 抽取和历史文件扫描的依赖。
注意:自 V2.0.11.1 起支持该功能。
2. 主要概念
IoTDB 订阅客户端包含 3 个核心概念:Topic、Consumer、Consumer Group,具体关系如下图

2.1 Topic
Topic 是 IoTDB 中可被订阅的数据空间。在表模型中,数据范围由以下配置共同确定:
database:数据库名称或匹配表达式。table:表名称或匹配表达式。column:列名称或匹配表达式。
普通订阅还可使用 start-time 和 end-time 限定事件时间(Event Time)范围。不同于 Kafka,IoTDB 可以在数据已经入库后再创建 Topic。
Topic 支持以下模式:
| 模式 | 说明 |
|---|---|
initial | 动态数据集。Consumer 可以持续消费符合条件的历史数据和后续新写入的数据。 |
snapshot | 静态数据集。以 Consumer Group 订阅 Topic 的时刻为边界生成快照,不持续包含快照之后的新数据。 |
incremental | 共识订阅。基于 IoTConsensus 共识写入日志消费实时增量数据,只消费 Consumer Group 首次成功订阅该 Topic 后进入共识链路的数据,不回溯历史数据。 |
表模型支持的输出格式如下:
| 格式 | 说明 |
|---|---|
SubscriptionRecordHandler | 按行消费数据。 |
SubscriptionTsFileHandler | 按 TsFile 消费数据。 |
共识订阅只支持
SubscriptionRecordHandler,不支持SubscriptionTsFileHandler。
2.2 Consumer
Consumer 是订阅客户端,负责订阅 Topic、接收数据并提交消费进度。
当前表模型提供 Pull Consumer。用户代码需要主动调用 poll() 拉取数据,并可选择自动或手动提交消费进度。
2.3 Consumer Group
拥有相同 Consumer Group ID 的 Consumers 属于同一个 Consumer Group:
- 一个 Consumer Group 可以包含多个 Consumers,一个 Consumer 只能加入一个 Consumer Group。
- 同一 Consumer Group 内,订阅相同 Topic 的 Consumers 共同分担消费负载;一条消息只分配给组内的一个 Consumer。
- 不同 Consumer Group 之间相互独立,可分别消费同一个 Topic。
- 未成功提交的消息在 Consumer 重启或重新分配后可能再次投递,因此数据订阅提供至少一次(At-least-once)语义,不提供精确一次(Exactly-once)语义。
2.4 注意事项
- 实时数据在单个 Region 内按照写入到达顺序消费,不保证跨 Region 的全局顺序,也不保证按照 Event Time 排序。
- Consumer 重启或重新分配后,未提交的消息可能再次投递。
- 下游系统应通过幂等写入或去重机制处理重复数据;依赖事件时间顺序时,需要自行处理乱序数据。
3. SQL 语句
3.1 Topic 管理
IoTDB 支持通过 SQL 语句对 Topic 进行创建、删除、查看操作。Topic状态变化如下图所示:

3.1.1 创建 Topic
SQL 语句为:
CREATE TOPIC [IF NOT EXISTS] <topicName>
WITH (
[<parameter> = <value>,],
);IF NOT EXISTS 语义:用于避免因 Topic 已经存在而报错。
3.1.1.1 普通订阅
普通订阅支持 initial 和 snapshot 模式,可通过数据库、表、列和时间范围筛选数据。
| 参数 | 默认值 | 说明 |
|---|---|---|
database | .* | 订阅的数据库,支持正则表达式。 |
table | .* | 订阅的表,支持正则表达式。 |
column | 全部列 | 订阅的列,支持正则表达式。 |
start-time | MIN_VALUE | Event Time 开始时间。支持 ISO 时间、与数据库时间戳精度一致的 long 值和特殊值 now。 |
end-time | MAX_VALUE | Event Time 结束时间,格式同 start-time。 |
mode | initial | 支持 initial和 snapshot。 initial模式导出全量+增量的数据,即导出全量后会持续导出增量数据。 snapshot模式仅导出当前时刻的快照数据,即当前时刻的全量数据。 |
format | SubscriptionRecordHandler | 支持 SubscriptionRecordHandler 和 SubscriptionTsFileHandler。 |
order-mode | leader-only | 支持 leader-only、multi-writer 和 per-writer。 |
loose-range | "" | 是否对数据范围和时间范围进行粗筛。 |
strict | true | 是否严格按照 Topic 范围筛选数据。 |
processor | do-nothing-processor | 对原始订阅数据应用的处理插件。 |
时间范围说明:
start-time=MIN_VALUE且end-time=now:只订阅历史数据。start-time=now且end-time=MAX_VALUE:只订阅实时数据。now表示 Topic 的创建时间。
示例:
- 全量订阅
CREATE TOPIC data_all;- 按数据库、表和时间范围订阅
CREATE TOPIC IF NOT EXISTS topic_table1
WITH (
'database' = 'database1',
'table' = 'table1',
'start-time' = '2024-11-01',
'end-time' = '2024-11-30',
'strict' = 'false'
);3.1.1.2 共识订阅
创建 Topic 时指定 mode=incremental 即可启用共识订阅。
- 共识订阅仅适用于使用 IoTConsensus 的 DataRegion。
- 共识订阅只面向实时增量数据。消费起点在 Consumer Group 首次成功订阅 Topic 时确定。
- 订阅前已经写入的数据不会补发,也不能通过时间范围回溯历史数据。
共识订阅支持以下配置:
| 参数 | 默认值 | 说明 |
|---|---|---|
database | .* | 数据库匹配范围。 |
table | .* | 表匹配范围。 |
column | 全部列 | 列匹配范围,支持正则表达式。 |
format | SubscriptionRecordHandler | 仅支持 SubscriptionRecordHandler。 |
retention.bytes | 536870912 | WAL 空间保留上限,单位为字节。支持正 long 或 -1;-1 表示不限制空间。 |
retention.ms | -1 | WAL 时间保留上限,单位为毫秒。支持正 long 或 -1;-1 表示不限制时间。 |
retention.bytes 和 retention.ms 不能设置为 0、小于 -1 或非 long。同时设置空间和时间限制时,任一限制达到上限后,旧 WAL 均可能被清理。已清理的 WAL 无法用于消费恢复或数据重放。Topic 创建后不能修改这两项配置。
示例:
CREATE TOPIC IF NOT EXISTS consensus_topic
WITH (
'mode' = 'incremental',
'database' = 'factory_db',
'table' = 'sensor_data',
'format' = 'SubscriptionRecordHandler',
'column' = '.*',
'retention.bytes' = '536870912',
'retention.ms' = '3600000'
);3.1.2 删除 Topic
Topic 在没有被订阅的情况下,才能被删除,Topic 被删除时,其相关的消费进度都会被清理。
DROP TOPIC [IF EXISTS] <topicName>;IF NOT EXISTS 语义:仅当指定 Topic 存在时执行删除,避免因 Topic 不存在而报错。
3.1.3 查看 Topic
SHOW TOPICS;
SHOW TOPIC <topicName>;结果集:
[TopicName|TopicConfigs]- TopicName:主题 ID
- TopicConfigs:主题配置,当前仅包含创建 Topic 时
WITH子句中设置的参数。
3.2 查看订阅状态
查看所有订阅关系:
SHOW SUBSCRIPTIONS;
SHOW SUBSCRIPTIONS ON <topicName>;结果集:
[SubscriptionID|TopicName|ConsumerGroupName|SubscribedConsumers]SubscriptionID:订阅关系的唯一标识,可用于管理或删除指定订阅关系。TopicName:Topic 名称。ConsumerGroupName:Consumer Group ID。SubscribedConsumers:当前订阅该 Topic 的 Consumer ID 集合。
4. API 接口
除 SQL 语句外,IoTDB 还支持通过 Java 原生接口使用数据订阅功能。详细语法参见数据订阅 API。
4.1 操作鉴权
元信息鉴权(SQL 语句权限)
- 创建、删除 Topic 需要
SYSTEM权限; - 查询 Topic 和订阅关系时,拥有
SYSTEM权限的用户可查看全局资源,普通用户仅能查看自身资源。
- 创建、删除 Topic 需要
运行时鉴权(客户端参数)
- Consumer 通过
Properties或 Builder 提供username和password。调用open()时校验身份,调用subscribe()时将鉴权信息传递给底层 subscription pipe,消费过程中则动态校验相关数据的查询权限;遇到无权访问的数据时,根据if-no-privileges等配置执行跳过或报错处理。
- Consumer 通过
Consumer Group 鉴权一致性限制
- 同一个 Consumer Group 下的所有 Consumers 必须使用相同的用户名和密码。第一个成功打开的 Consumer 确定该组的鉴权基准;后续 Consumer 的鉴权信息一致时允许加入,不一致时拒绝加入。