数据订阅API
数据订阅API
IoTDB 表模型数据订阅 API 允许应用通过 Java SDK 管理 Topic、主动拉取订阅数据并提交消费进度。功能定义和 SQL 语法参见:数据订阅。
注意:自 V2.0.11.1 起支持该功能。
1. 核心步骤
- 创建 Topic:使用
ISubscriptionTableSession创建 Topic,通过数据库、表和时间范围定义订阅数据。 - 订阅 Topic:Consumer 只能订阅已经创建的 Topic;同一个 Consumer Group 内订阅相同 Topic 的 Consumers 共同分担数据。
- 消费数据:调用
poll()主动拉取SubscriptionMessage。 - 提交进度:使用默认的自动提交,或关闭自动提交后调用
commitSync()、commitAsync()。 - 取消订阅:调用
unsubscribe();Consumer 关闭时会退出 Consumer Group。
对于 mode=consensus 的 Topic,需要注意:
- 共识订阅仅适用于使用 IoTConsensus 的 DataRegion。
- 不支持通过
start-time或end-time指定消费范围。 - 消费起点在 Consumer Group 首次成功订阅 Topic 时确定,订阅前数据不会补发。
- 只支持
SubscriptionRecordHandler,不支持 TsFile 格式。
不同客户端即使配置了相同的 Consumer Group ID 和 Consumer ID,服务端仍会将其视为两个独立连接,由两个客户端共同分担消费负载。
2. 详细步骤
本章节用于说明开发的核心流程,并未演示所有的参数和接口,如需了解全部功能及参数请参见参数与接口。
2.1 创建maven项目
创建一个maven项目,并导入以下依赖(JDK >= 17, Maven >= 3.6)
<dependencies>
<dependency>
<groupId>org.apache.iotdb</groupId>
<artifactId>iotdb-session</artifactId>
<!-- 版本号与数据库版本号相同 -->
<version>${project.version}</version>
</dependency>
</dependencies>2.2 使用示例
2.2.1 管理普通 Topic
import java.util.Properties;
import org.apache.iotdb.rpc.subscription.config.TopicConstant;
import org.apache.iotdb.session.subscription.ISubscriptionTableSession;
import org.apache.iotdb.session.subscription.SubscriptionTableSessionBuilder;
public class TopicOperationExample {
public static void main(String[] args) throws Exception {
try (final ISubscriptionTableSession session =
new SubscriptionTableSessionBuilder()
.host("127.0.0.1")
.port(6667)
.username("root")
.password("TimechoDB@2021")
.build()) {
final Properties config = new Properties();
config.put(TopicConstant.DATABASE_KEY, "db.*");
config.put(TopicConstant.TABLE_KEY, "test.*");
config.put(TopicConstant.START_TIME_KEY, 25);
config.put(TopicConstant.END_TIME_KEY, 75);
config.put(TopicConstant.STRICT_KEY, "true");
config.put(
TopicConstant.FORMAT_KEY,
TopicConstant.FORMAT_RECORD_HANDLER_VALUE);
session.createTopicIfNotExists("topic1", config);
session.getTopic("topic1").ifPresent(System.out::println);
session.getSubscriptions("topic1").forEach(System.out::println);
}
}
}2.2.2 创建共识订阅 Topic
以下示例使用 Topic 配置字符串创建共识订阅 Topic:
import java.util.Properties;
import org.apache.iotdb.session.subscription.ISubscriptionTableSession;
import org.apache.iotdb.session.subscription.SubscriptionTableSessionBuilder;
public class ConsensusTopicOperationExample {
public static void main(String[] args) throws Exception {
try (final ISubscriptionTableSession session =
new SubscriptionTableSessionBuilder()
.host("127.0.0.1")
.port(6667)
.username("root")
.password("TimechoDB@2021")
.build()) {
final Properties config = new Properties();
config.setProperty("mode", "consensus");
config.setProperty("database", "factory_db");
config.setProperty("table", "sensor_data");
config.setProperty("column", ".*");
config.setProperty("format", "SubscriptionRecordHandler");
config.setProperty("retention.bytes", "536870912");
config.setProperty("retention.ms", "3600000");
session.createTopicIfNotExists("consensus_topic", config);
session.getTopic("consensus_topic").ifPresent(System.out::println);
}
}
}创建共识订阅 Topic 时不要设置
start-time、end-time、strict、order-mode等不受支持的配置。
2.2.3 按行消费数据
行数据通过 SubscriptionMessage.getResultSets() 获取。
import java.util.List;
import org.apache.iotdb.session.subscription.consumer.ISubscriptionTablePullConsumer;
import org.apache.iotdb.session.subscription.consumer.table.SubscriptionTablePullConsumerBuilder;
import org.apache.iotdb.session.subscription.payload.SubscriptionMessage;
import org.apache.iotdb.session.subscription.payload.SubscriptionRecordHandler;
import org.apache.tsfile.read.query.dataset.ResultSet;
public class RecordSubscriptionExample {
public static void main(String[] args) throws Exception {
try (final ISubscriptionTablePullConsumer consumer =
new SubscriptionTablePullConsumerBuilder()
.host("127.0.0.1")
.port(6667)
.username("root")
.password("TimechoDB@2021")
.consumerId("c1")
.consumerGroupId("cg1")
.build()) {
consumer.open();
consumer.subscribe("topic1");
while (true) {
final List<SubscriptionMessage> messages = consumer.poll(10_000L);
for (final SubscriptionMessage message : messages) {
for (final ResultSet resultSet : message.getResultSets()) {
final SubscriptionRecordHandler.SubscriptionResultSet recordSet =
(SubscriptionRecordHandler.SubscriptionResultSet) resultSet;
System.out.println(recordSet.getDatabaseName());
System.out.println(recordSet.getTableName());
System.out.println(recordSet.getColumnNames());
System.out.println(recordSet.getColumnTypes());
System.out.println(recordSet.getColumnCategories());
while (recordSet.hasNext()) {
System.out.println(recordSet.nextRecord());
}
}
}
// autoCommit 默认为 true。
}
}
}
}该消费代码也适用于共识订阅 Topic。共识订阅返回的消息只包含 Consumer Group 建立订阅关系后进入共识链路的数据。
2.2.4 手动提交消费进度
需要在业务处理成功后再提交进度时,可关闭自动提交:
import java.util.List;
import org.apache.iotdb.session.subscription.consumer.ISubscriptionTablePullConsumer;
import org.apache.iotdb.session.subscription.consumer.table.SubscriptionTablePullConsumerBuilder;
import org.apache.iotdb.session.subscription.payload.SubscriptionMessage;
import org.apache.iotdb.session.subscription.payload.SubscriptionRecordHandler;
import org.apache.tsfile.read.query.dataset.ResultSet;
public class ManualCommitSubscriptionExample {
public static void main(String[] args) throws Exception {
try (final ISubscriptionTablePullConsumer consumer =
new SubscriptionTablePullConsumerBuilder()
.host("127.0.0.1")
.port(6667)
.username("root")
.password("TimechoDB@2021")
.consumerId("c1")
.consumerGroupId("cg1")
.autoCommit(false)
.build()) {
consumer.open();
consumer.subscribe("consensus_topic");
while (true) {
final List<SubscriptionMessage> messages = consumer.poll(10_000L);
if (messages.isEmpty()) {
continue;
}
for (final SubscriptionMessage message : messages) {
for (final ResultSet resultSet : message.getResultSets()) {
final SubscriptionRecordHandler.SubscriptionResultSet recordSet =
(SubscriptionRecordHandler.SubscriptionResultSet) resultSet;
while (recordSet.hasNext()) {
// 业务处理应具备幂等性。
System.out.println(recordSet.nextRecord());
}
}
}
consumer.commitSync(messages);
}
}
}
}如果业务处理完成后、提交成功前 Consumer 异常退出,未确认消息会在重启或重新分配后再次投递。因此数据订阅提供至少一次(At-least-once)语义,不提供精确一次(Exactly-once)语义。
客户端处理器仍缓存消息时,可在关闭前调用 drainBufferedMessages(),处理并提交返回的消息。
2.2.5 订阅 TsFile
本场景仅适用于普通订阅,不适用于
mode=consensus的 Topic。
创建 TsFile 格式 Topic:
CREATE TOPIC topic_tsfile
WITH (
'database' = 'database1',
'table' = 'table1',
'format' = 'SubscriptionTsFileHandler'
);消费 TsFile:
import org.apache.iotdb.session.subscription.consumer.ISubscriptionTablePullConsumer;
import org.apache.iotdb.session.subscription.consumer.table.SubscriptionTablePullConsumerBuilder;
import org.apache.iotdb.session.subscription.payload.SubscriptionMessage;
import org.apache.iotdb.session.subscription.payload.SubscriptionTsFileHandler;
import org.apache.tsfile.read.v4.ITsFileReader;
public class TsFileSubscriptionExample {
public static void main(String[] args) throws Exception {
try (final ISubscriptionTablePullConsumer consumer =
new SubscriptionTablePullConsumerBuilder()
.host("127.0.0.1")
.port(6667)
.username("root")
.password("TimechoDB@2021")
.consumerId("c1")
.consumerGroupId("cg1")
.build()) {
consumer.open();
consumer.subscribe("topic_tsfile");
while (true) {
for (final SubscriptionMessage message : consumer.poll(10_000L)) {
final SubscriptionTsFileHandler handler = message.getTsFile();
try (final ITsFileReader reader = handler.openTableReader()) {
// 读取订阅到的表模型 TsFile。
}
}
}
}
}
}3. 参数与接口
3.1 常用参数
3.1.1 Consumer 公共配置
| 参数 | 默认值 | 说明 |
|---|---|---|
host | 127.0.0.1 | DataNode RPC Host。 |
port | 6667 | DataNode RPC Port。 |
nodeUrls | 无 | 多个 DataNode RPC 地址。 |
username | root | 用户名。 |
password | root | 密码。 |
encryptedPassword | 无 | 加密密码。 |
consumerId | 自动分配 | Consumer 打开后自动分配全局唯一 ID。 |
consumerGroupId | 自动分配 | Consumer 打开后自动分配全局唯一 ID。 |
heartbeatIntervalMs | 30000,最小 1000 | 心跳间隔,单位为毫秒。 |
endpointsSyncIntervalMs | 120000,最小 5000 | 集群端点同步间隔,单位为毫秒。 |
fileSaveDir | <user.dir>/iotdb-subscription | TsFile 临时保存目录。 |
fileSaveFsync | false | 保存 TsFile 时是否执行 fsync。 |
connectionTimeoutInMs | 0 | 连接超时时间。 |
maxPollParallelism | 1 | 最大并行 Poll 数。 |
3.1.2 Pull Consumer 特殊配置
| 参数 | 默认值 | 说明 |
|---|---|---|
autoCommit | true | 是否自动提交消费进度;为 false 时需要手动调用提交接口。 |
autoCommitIntervalMs | 5000,最小 500 | 自动提交间隔,仅在 autoCommit=true 时生效。 |
3.2 接口介绍
3.2.1 ISubscriptionTableSession
用于管理表模型 Topic,并查询 Topic 和订阅关系。
| 方法 | 说明 | 返回值 | 主要异常 |
|---|---|---|---|
open() | 打开 Session 并连接 IoTDB。 | void | IoTDBConnectionException |
createTopic(String topicName) | 创建使用默认配置的 Topic。 | void | IoTDBConnectionException、StatementExecutionException |
createTopicIfNotExists(String topicName) | Topic 不存在时创建。 | void | 同上 |
createTopic(String topicName, Properties properties) | 使用属性创建 Topic。 | void | 同上 |
createTopicIfNotExists(String topicName, Properties properties) | Topic 不存在时使用属性创建。 | void | 同上 |
alterTopic(String topicName, Properties properties) | 修改 Topic 的可变配置。 | void | 同上 |
alterTopicOwner(String topicName, String ownerId, long ownerEpoch) | 修改 Topic Owner 信息。 | void | 同上 |
alterTopicOwner(String topicName, String ownerId, long ownerEpoch, Long maxOwnerEpoch) | 修改 Topic Owner 及最大 Epoch。 | void | 同上 |
dropTopic(String topicName) | 删除 Topic。 | void | 同上 |
dropTopicIfExists(String topicName) | Topic 存在时删除。 | void | 同上 |
dropSubscription(String subscriptionId) | 删除指定订阅关系。 | void | 同上 |
dropSubscriptionIfExists(String subscriptionId) | 订阅关系存在时删除。 | void | 同上 |
getTopics() | 获取全部 Topic。 | Set<Topic> | 同上 |
getTopic(String topicName) | 获取指定 Topic,不存在时为空。 | Optional<Topic> | 同上 |
getSubscriptions() | 获取全部订阅关系。 | Set<Subscription> | 同上 |
getSubscriptions(String topicName) | 获取指定 Topic 的订阅关系。 | Set<Subscription> | 同上 |
close() | 关闭 Session。 | void | Exception |
3.2.2 ISubscriptionTablePullConsumer
用于订阅 Topic、主动拉取消息并提交消费进度。
| 方法 | 说明 | 返回值 | 主要异常 |
|---|---|---|---|
open() | 打开 Consumer。 | void | SubscriptionException |
subscribe(String topicName) | 订阅单个 Topic。 | void | SubscriptionException |
subscribe(String... topicNames) | 订阅多个 Topic。 | void | SubscriptionException |
subscribe(Set<String> topicNames) | 订阅多个 Topic。 | void | SubscriptionException |
unsubscribe(String topicName) | 取消订阅单个 Topic。 | void | SubscriptionException |
unsubscribe(String... topicNames) | 取消订阅多个 Topic。 | void | SubscriptionException |
unsubscribe(Set<String> topicNames) | 取消订阅多个 Topic。 | void | SubscriptionException |
poll(Duration timeout) | 拉取消息,参数为无消息时的最长等待时间,不限制返回数量。 | List<SubscriptionMessage> | SubscriptionException |
poll(long timeoutMs) | 拉取消息,参数为最长等待毫秒数,不限制返回数量。 | List<SubscriptionMessage> | SubscriptionException |
poll(Set<String> topicNames, Duration timeout) | 从指定 Topic 集合拉取消息。 | List<SubscriptionMessage> | SubscriptionException |
poll(Set<String> topicNames, long timeoutMs) | 从指定 Topic 集合拉取消息。 | List<SubscriptionMessage> | 无 |
drainBufferedMessages() | 排空客户端处理器中缓存的消息。 | List<SubscriptionMessage> | SubscriptionException |
commitSync(SubscriptionMessage message) | 同步提交单条消息。 | void | SubscriptionException |
commitSync(Iterable<SubscriptionMessage> messages) | 同步提交多条消息。 | void | SubscriptionException |
commitAsync(SubscriptionMessage message) | 异步提交单条消息。 | CompletableFuture<Void> | 无 |
commitAsync(Iterable<SubscriptionMessage> messages) | 异步提交多条消息。 | CompletableFuture<Void> | 无 |
commitAsync(SubscriptionMessage message, AsyncCommitCallback callback) | 异步提交单条消息并回调结果。 | void | 无 |
commitAsync(Iterable<SubscriptionMessage> messages, AsyncCommitCallback callback) | 异步提交多条消息并回调结果。 | void | 无 |
seekToBeginning(String topicName) | 将消费位置移动到开始位置。 | void | SubscriptionException |
seekToEnd(String topicName) | 将消费位置移动到结束位置。 | void | SubscriptionException |
positions(String topicName) | 获取当前位置。 | TopicProgress | SubscriptionException |
committedPositions(String topicName) | 获取已提交位置。 | TopicProgress | SubscriptionException |
seek(String topicName, TopicProgress topicProgress) | 移动到指定消费位置。 | void | SubscriptionException |
seekAfter(String topicName, TopicProgress topicProgress) | 移动到指定消费位置之后。 | void | SubscriptionException |
getConsumerId() | 获取 Consumer ID。 | String | 无 |
getConsumerGroupId() | 获取 Consumer Group ID。 | String | 无 |
allTopicMessagesHaveBeenConsumed() | 判断所有 Topic 消息是否已消费完成。 | boolean | 无 |
close() | 关闭 Consumer,退出 Consumer Group。 | void | Exception |
3.2.3 SubscriptionMessage
SubscriptionMessage 是 poll() 返回的基本消息单元。数据形式由 Topic 的 format 决定:按行格式通过 getResultSets() 读取,TsFile 格式通过 getTsFile() 读取。
| 方法 | 说明 | 返回值 |
|---|---|---|
getMessageType() | 获取消息类型。 | short |
isTimeSelected() | 判断消息是否经过时间范围筛选。 | boolean |
getResultSets() | 获取按行格式结果集。 | List<ResultSet> |
getRecordTabletIterator() | 以 Tablet 迭代器读取按行数据。 | Iterator<Tablet> |
getTsFile() | 获取 TsFile 处理对象。 | SubscriptionTsFileHandler |
getWatermarkTimestamp() | 获取 Watermark 消息携带的时间戳。 | long |
estimateSize() | 估算消息占用的堆内存字节数。 | long |
removeUserData() | 释放消息中的用户数据。 | void |
调用与消息格式不兼容的读取方法时,会抛出 SubscriptionIncompatibleHandlerException。消息数据已释放后再次读取,可能抛出 SubscriptionRuntimeException。
3.2.4 SubscriptionRecordHandler.SubscriptionResultSet
表模型按行数据的结果集实现。
| 方法 | 说明 | 返回值 |
|---|---|---|
getDatabaseName() | 获取数据库名。 | String |
getTableName() | 获取表名。 | String |
getColumnNames() | 获取列名。 | 列名列表 |
getColumnTypes() | 获取列数据类型。 | 列类型列表 |
getColumnCategories() | 获取列类别。 | 列类别列表 |
hasNext() | 判断是否还有下一行。 | boolean |
nextRecord() | 读取下一行。 | RowRecord |
3.2.5 SubscriptionTsFileHandler
用于处理普通订阅传输到客户端的 TsFile。共识订阅不支持该 Handler。
| 方法 | 说明 | 返回值 |
|---|---|---|
getDatabaseName() | 获取 TsFile 所属数据库。 | String |
openTableReader() | 打开表模型 TsFile Reader。 | ITsFileReader |
文件不是表模型 TsFile 时,openTableReader() 会抛出 SubscriptionIncompatibleHandlerException;读取文件还可能抛出 IOException。