数据订阅
数据订阅
1. 功能介绍
IoTDB 数据订阅模块(又称 IoTDB 订阅客户端)提供了一种区别于定时查询的流式数据消费方式。它参考 Kafka 等消息队列产品的基本概念和逻辑,为应用提供 Topic 管理、数据订阅和消费进度提交能力。
数据订阅并不是为了完全替代消息队列,主要适用于以下场景:
- 持续获取最新数据:持续拉取新写入的数据,无需频繁执行定时查询,适用于大屏展示、组态监控等场景。
- 简化第三方系统对接:Flink、Kafka、DataX、Camel、MySQL、PostgreSQL 等系统可以作为订阅客户端主动获取数据,无需在 IoTDB 内为每个系统分别开发推送组件。
- 增量备份 TsFile:订阅新生成的 TsFile,将增量数据归档到指定存储位置。
注意:自 V2.0.11.1 起支持该功能。
2. 主要概念
IoTDB 订阅客户端包含 Topic、Consumer 和 Consumer Group 三个核心概念。

2.1 Topic
Topic 是 IoTDB 中可被订阅的数据空间。在树模型中,数据范围由路径或路径模式确定,例如 root.**;普通订阅还可使用 start-time 和 end-time 限定事件时间(Event Time)范围。
不同于 Kafka,IoTDB 可以在数据已经入库后再创建 Topic。Topic 可以按行或按 TsFile 输出数据。
树模型支持以下订阅模式:
| 模式 | 说明 |
|---|---|
initial | 动态数据集。Consumer 可以持续消费符合条件的历史数据和后续新写入的数据。 |
snapshot | 静态数据集。以 Consumer Group 订阅 Topic 的时刻为边界生成快照,不持续包含快照之后的新数据。 |
incremental | 基于共识进度的实时增量数据集。只消费 Consumer Group 订阅 Topic 后新写入的数据,不回溯订阅前的历史数据。 |
树模型支持以下输出格式:
| 格式 | 说明 |
|---|---|
SubscriptionRecordHandler | 以 SubscriptionMessage 中的 ResultSet 或 Tablet 形式消费数据。 |
SubscriptionTsFileHandler | 使用 SubscriptionTsFileHandler 消费 TsFile。 |
2.2 Consumer
Consumer 是订阅客户端,负责接收和处理 Topic 数据。树模型提供两种 Consumer:
ISubscriptionTreePullConsumer:Pull 模式,通过SubscriptionTreePullConsumerBuilder构建,由用户代码主动调用poll()获取数据。ISubscriptionTreePushConsumer:Push 模式,通过SubscriptionTreePushConsumerBuilder构建,由新到达的数据触发用户实现的ConsumeListener。
2.3 Consumer Group
拥有相同 Consumer Group ID 的 Consumers 属于同一个 Consumer Group:
- 一个 Consumer Group 可以包含多个 Consumers,一个 Consumer 只能加入一个 Consumer Group。
- 一个 Consumer Group 中可以包含 Pull Consumer 和 Push Consumer。
- 一个 Topic 不要求被 Consumer Group 中的所有 Consumers 订阅。
- 同一 Consumer Group 内,订阅相同 Topic 的 Consumers 共同分担消费负载,一条消息只分配给组内的一个 Consumer。
- 不同 Consumer Group 相互独立,可以分别消费同一个 Topic。
- 未成功提交的消息在 Consumer 重启或重新分配后可能再次投递,因此数据订阅提供至少一次(At-least-once)语义,不提供精确一次(Exactly-once)语义。
3. SQL 语句
3.1 Topic 管理

3.1.1 创建 Topic
CREATE TOPIC [IF NOT EXISTS] <topicName>
WITH (
[<parameter> = <value>,]
);IF NOT EXISTS 用于避免因 Topic 已经存在而报错。
- Topic 配置
| 参数 | 默认值 | 说明 |
|---|---|---|
path | root.** | Topic 对应的时间序列路径或路径模式。 |
start-time | MIN_VALUE | Event Time 开始时间。支持 ISO 时间、与数据库时间戳精度一致的 long 值和特殊值 now。 |
end-time | MAX_VALUE | Event Time 结束时间,格式同 start-time。 |
processor | do-nothing-processor | 对原始订阅数据应用的处理插件。 |
format | SubscriptionRecordHandler | 支持 SubscriptionRecordHandler 和 SubscriptionTsFileHandler。 |
mode | initial | 支持 initial、snapshot 和 incremental。initial模式导出全量+增量的数据,即导出全量后会持续导出增量数据。snapshot模式仅导出当前时刻的快照数据,即当前时刻的全量数据。incremental 模式仅导出启动订阅后产生的增量数据。可通过 retention.bytes 和 retention.ms 分别设置保留数据量上限(字节)和保留时间上限(毫秒)。该模式依赖 DataRegion 使用 IoTConsensus。 |
loose-range | "" | 是否对路径和时间范围进行粗筛。
|
时间范围示例:
start-time=MIN_VALUE且end-time=now:只订阅历史数据。start-time=now且end-time=MAX_VALUE:只订阅实时数据。now表示 Topic 的创建时间。
创建示例
全量订阅:
CREATE TOPIC root_all;订阅指定路径和时间范围:
CREATE TOPIC IF NOT EXISTS db_timerange
WITH (
'path' = 'root.db.**',
'start-time' = '2023-01-01',
'end-time' = '2023-12-31'
);订阅新增 TsFile:
CREATE TOPIC topic_all_tsfile
WITH (
'path' = 'root.**',
'format' = 'SubscriptionTsFileHandler'
);创建基于共识进度的实时增量 Topic:
CREATE TOPIC topic_realtime_consensus
WITH (
'path' = 'root.db.**',
'mode' = 'incremental',
'retention.bytes' = '536870912',
'retention.ms' = '3600000'
);3.1.2 删除 Topic
只有未被订阅的 Topic 才能被删除。Topic 删除后,其相关消费进度会被清理。
DROP TOPIC [IF EXISTS] <topicName>;IF EXISTS 表示仅当 Topic 存在时执行删除,避免因 Topic 不存在而报错。
3.1.3 查看 Topic
SHOW TOPICS;
SHOW TOPIC <topicName>;结果集:
[TopicName|TopicConfigs]TopicName:Topic 名称。TopicConfigs:创建 Topic 时WITH子句中设置的配置。
查询不存在的 Topic 时返回空结果集。
3.2 查看订阅状态
SHOW SUBSCRIPTIONS;
SHOW SUBSCRIPTIONS ON <topicName>;结果集:
[SubscriptionID|TopicName|ConsumerGroupName|SubscribedConsumers]SubscriptionID:订阅关系 ID。TopicName:Topic 名称。ConsumerGroupName:Consumer Group ID。SubscribedConsumers:该 Consumer Group 中订阅此 Topic 的所有 Consumer ID。
4. API 接口
除 SQL 语句外,IoTDB 还支持通过 Java 原生接口使用数据订阅功能。详细语法参见页面:数据订阅 API。
4.1 操作鉴权
元信息鉴权
- 创建、删除 Topic 需要
SYSTEM权限。 - 查询 Topic 和订阅关系时,拥有
SYSTEM权限的用户可以查看全局资源,普通用户只能查看自身资源。
- 创建、删除 Topic 需要
运行时鉴权
- Consumer 通过 Properties 或 Builder 提供
username和password:调用open()时校验身份。调用subscribe()时校验订阅操作。消费过程中动态校验相关数据的查询权限;遇到无权访问的数据时,根据if-no-privileges等配置跳过或报错。
- Consumer 通过 Properties 或 Builder 提供
Consumer Group 鉴权一致性限制
- 同一个 Consumer Group 下的所有 Consumers 必须使用相同的用户名和密码。第一个成功打开的 Consumer 确定该组的鉴权基准;后续 Consumer 的鉴权信息一致时允许加入,不一致时拒绝加入。
5. 常见问题
5.1 IoTDB 数据订阅与 Kafka 的区别是什么?
消费有序性
- Kafka 保证消息在单个 partition 内是有序的,当某个 topic 仅对应一个 partition 且只有一个 consumer 订阅了这个 topic,即可保证该 consumer(单线程) 消费该 topic 数据的顺序即为数据写入的顺序。
- IoTDB 订阅客户端不保证 consumer 消费数据的顺序即为数据写入的顺序,但会尽量反映数据写入的顺序。
消息送达语义
- Kafka 可以通过配置实现 Producer 和 Consumer 的 Exactly once 语义。
- IoTDB 订阅客户端目前无法提供 Consumer 的 Exactly once 语义。
5.2 使用数据订阅时需要注意什么?
- 实时数据在单个 Region 内按照写入到达顺序消费,不保证跨 Region 的全局顺序,也不保证按照 Event Time 排序。
- Consumer 重启或重新分配后,未提交的消息可能再次投递。
- 下游系统应通过幂等写入或去重机制处理重复数据;依赖事件时间顺序时,需要自行处理乱序数据。