数据订阅API
数据订阅API
IoTDB 树模型数据订阅 API 允许应用通过 Java SDK 管理 Topic、使用 Pull 或 Push Consumer 获取数据,并提交消费进度。功能定义和 SQL 语法参见:数据订阅。
注意:自 V2.0.11.1 起支持该功能。
1. 核心步骤
- 创建 Topic:通过路径和可选时间范围定义希望订阅的时间序列。
- 订阅 Topic:Consumer 只能订阅已经创建的 Topic。同一 Consumer Group 内订阅相同 Topic 的 Consumers 共同分担数据。
- 消费数据:Pull Consumer 调用
poll()主动拉取;Push Consumer通过ConsumeListener接收回调。 - 提交进度:Pull Consumer 可使用自动提交,也可关闭自动提交后调用
commitSync()、commitAsync()。 - 取消订阅:调用
unsubscribe();Consumer 关闭时会退出 Consumer Group 并取消该 Consumer 的现有订阅。
不同客户端即使配置了相同的 Consumer Group ID 和 Consumer ID,服务端仍会将其视为两个独立连接,并由两个客户端共同分担消费负载。
2. 详细步骤
2.1 创建 Maven 项目
创建 Maven 项目并引入 iotdb-session。依赖版本应与数据库版本保持一致。
<dependencies>
<dependency>
<groupId>org.apache.iotdb</groupId>
<artifactId>iotdb-session</artifactId>
<version>${project.version}</version>
</dependency>
</dependencies>运行环境要求:
- JDK 17 或更高版本。。
- Maven 3.6 或更高版本。
- 不要使用高版本客户端连接低版本服务端。
2.2 代码案例
2.2.1 Topic 操作
使用 SubscriptionTreeSessionBuilder 创建树模型订阅 Session。build() 只构建对象,执行 Topic 操作前仍需调用 open()。
import java.util.Optional;
import java.util.Properties;
import java.util.Set;
import org.apache.iotdb.rpc.subscription.config.TopicConstant;
import org.apache.iotdb.session.subscription.ISubscriptionTreeSession;
import org.apache.iotdb.session.subscription.SubscriptionTreeSessionBuilder;
import org.apache.iotdb.session.subscription.model.Topic;
public class TopicOperationExample {
public static void main(String[] args) throws Exception {
try (ISubscriptionTreeSession session =
new SubscriptionTreeSessionBuilder()
.host("127.0.0.1")
.port(6667)
.username("root")
.password("TimechoDB@2021")
.build()) {
session.open();
final Properties topicConfig = new Properties();
topicConfig.setProperty(TopicConstant.PATH_KEY, "root.**");
session.createTopicIfNotExists("allData", topicConfig);
final Set<Topic> topics = session.getTopics();
System.out.println(topics);
final Optional<Topic> allData = session.getTopic("allData");
allData.ifPresent(System.out::println);
}
}
}V2.0.6.x 之前的默认密码为
root。请以部署版本的实际账号配置为准,生产代码中不要硬编码密码。
2.2.2 Pull 模式消费 Record
SubscriptionMessage 可通过 getRecordTabletIterator() 返回 Tablet 迭代器。每个 Tablet 包含设备、时间戳、测点 Schema 和值。
import java.util.Iterator;
import java.util.List;
import org.apache.iotdb.session.subscription.consumer.ISubscriptionTreePullConsumer;
import org.apache.iotdb.session.subscription.consumer.tree.SubscriptionTreePullConsumerBuilder;
import org.apache.iotdb.session.subscription.payload.SubscriptionMessage;
import org.apache.iotdb.session.subscription.payload.SubscriptionMessageType;
import org.apache.tsfile.write.record.Tablet;
public class RecordSubscriptionExample {
public static void main(String[] args) throws Exception {
try (ISubscriptionTreePullConsumer consumer =
new SubscriptionTreePullConsumerBuilder()
.host("127.0.0.1")
.port(6667)
.username("root")
.password("TimechoDB@2021")
.consumerId("c1")
.consumerGroupId("cg1")
.build()) {
consumer.open();
consumer.subscribe("topic_all");
while (true) {
final List<SubscriptionMessage> messages = consumer.poll(10_000L);
for (final SubscriptionMessage message : messages) {
if (message.getMessageType()
!= SubscriptionMessageType.RECORD_HANDLER.getType()) {
continue;
}
final Iterator<Tablet> tablets = message.getRecordTabletIterator();
while (tablets.hasNext()) {
final Tablet tablet = tablets.next();
for (int row = 0; row < tablet.getRowSize(); row++) {
System.out.printf(
"device=%s, time=%d%n",
tablet.getDeviceId(), tablet.getTimestamp(row));
for (int column = 0; column < tablet.getSchemas().size(); column++) {
System.out.printf(
" %s=%s%n",
tablet.getSchemas().get(column).getMeasurementName(),
tablet.getValue(row, column));
}
}
}
}
}
}
}
}也可以使用 message.getResultSets() 获取 List<org.apache.tsfile.read.query.dataset.ResultSet>。
poll(timeoutMs) 中的参数表示没有可用消息时愿意等待的最长时间,不表示单次最多返回的消息数量。
2.2.3 手动提交消费进度
需要在业务处理成功后再提交消费进度时,可通过 Builder 关闭自动提交:
try (ISubscriptionTreePullConsumer consumer =
new SubscriptionTreePullConsumerBuilder()
.host("127.0.0.1")
.port(6667)
.username("root")
.password("TimechoDB@2021")
.consumerId("c1")
.consumerGroupId("cg1")
.autoCommit(false)
.build()) {
consumer.open();
consumer.subscribe("topic_all");
while (true) {
final List<SubscriptionMessage> messages = consumer.poll(10_000L);
if (messages.isEmpty()) {
continue;
}
for (final SubscriptionMessage message : messages) {
// 下游处理应具备幂等性。
System.out.println(message);
}
// 只在整个批次处理成功后提交。
consumer.commitSync(messages);
}
}如果业务处理完成后、提交成功前 Consumer 异常退出,未确认消息会在重启或重新分配后再次投递。因此数据订阅提供至少一次(At-least-once)语义,不提供精确一次(Exactly-once)语义。
2.2.4 Push 模式消费
Push Consumer 通过 SubscriptionTreePushConsumerBuilder 配置消费回调和确认策略:
AckStrategy.BEFORE_CONSUME:在调用消费回调前确认进度。进程在确认后、处理前退出时,可能导致消息未被业务处理。AckStrategy.AFTER_CONSUME:在消费回调成功后确认进度。进程在处理后、确认前退出时,消息可能再次投递。
import org.apache.iotdb.session.subscription.consumer.AckStrategy;
import org.apache.iotdb.session.subscription.consumer.ConsumeResult;
import org.apache.iotdb.session.subscription.consumer.ISubscriptionTreePushConsumer;
import org.apache.iotdb.session.subscription.consumer.tree.SubscriptionTreePushConsumerBuilder;
public class PushSubscriptionExample {
public static void main(String[] args) throws Exception {
try (ISubscriptionTreePushConsumer consumer =
new SubscriptionTreePushConsumerBuilder()
.host("127.0.0.1")
.port(6667)
.username("root")
.password("TimechoDB@2021")
.consumerId("c2")
.consumerGroupId("cg1")
.ackStrategy(AckStrategy.AFTER_CONSUME)
.autoPollIntervalMs(100L)
.autoPollTimeoutMs(10_000L)
.consumeListener(
message -> {
try {
System.out.println(message);
return ConsumeResult.SUCCESS;
} catch (Exception e) {
return ConsumeResult.FAILURE;
}
})
.build()) {
consumer.open();
consumer.subscribe("topic_all");
// 示例程序保持运行;实际应用应由自身生命周期管理组件控制。
Thread.currentThread().join();
}
}
}Push Consumer 会按照 autoPollIntervalMs 和 autoPollTimeoutMs 自动拉取数据并调用监听器。
2.2.5 订阅 TsFile
首先创建 TsFile 格式 Topic。格式值是 SubscriptionTsFileHandler:
CREATE TOPIC topic_all_tsfile
WITH (
'path' = 'root.**',
'format' = 'SubscriptionTsFileHandler'
);然后使用 Pull Consumer 获取文件。getTsFile() 返回 SubscriptionTsFileHandler,树模型文件应通过 openTreeReader() 打开:
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.util.List;
import org.apache.iotdb.session.subscription.consumer.ISubscriptionTreePullConsumer;
import org.apache.iotdb.session.subscription.consumer.tree.SubscriptionTreePullConsumerBuilder;
import org.apache.iotdb.session.subscription.payload.SubscriptionMessage;
import org.apache.iotdb.session.subscription.payload.SubscriptionMessageType;
import org.apache.iotdb.session.subscription.payload.SubscriptionTsFileHandler;
import org.apache.tsfile.read.v4.ITsFileTreeReader;
public class TsFileSubscriptionExample {
public static void main(String[] args) throws Exception {
try (ISubscriptionTreePullConsumer consumer =
new SubscriptionTreePullConsumerBuilder()
.host("127.0.0.1")
.port(6667)
.username("root")
.password("TimechoDB@2021")
.consumerId("c1")
.consumerGroupId("cg1")
.fileSaveDir("/Users/iotdb/Downloads/subscription-cache")
.autoCommit(false)
.build()) {
consumer.open();
consumer.subscribe("topic_all_tsfile");
while (true) {
final List<SubscriptionMessage> messages = consumer.poll(10_000L);
for (final SubscriptionMessage message : messages) {
if (message.getMessageType() != SubscriptionMessageType.TS_FILE.getType()) {
continue;
}
final SubscriptionTsFileHandler handler = message.getTsFile();
try (ITsFileTreeReader reader = handler.openTreeReader()) {
System.out.println(reader.getAllDeviceIds());
}
final Path archiveDir = Paths.get("/Users/iotdb/Downloads/archive");
Files.createDirectories(archiveDir);
final Path target = archiveDir.resolve(handler.getFile().getName());
handler.copyFile(target);
}
if (!messages.isEmpty()) {
consumer.commitSync(messages);
}
}
}
}
}3. 常用接口说明
3.1 参数列表
3.1.1 Consumer 公共配置
| Builder 参数 | 默认值 | 说明 |
|---|---|---|
host | 127.0.0.1 | DataNode RPC Host。 |
port | 6667 | DataNode RPC Port。 |
nodeUrls | 127.0.0.1:6667 | DataNode RPC 地址列表;与 host/port 同时填写时取并集。 |
username | root | 用户名。 |
password | TimechoDB@2021 | 密码;V2.0.6.x 之前默认值为 root。 |
encryptedPassword | 无 | 已加密的密码。 |
consumerGroupId | 自动分配 | Consumer Group ID。 |
consumerId | 自动分配 | Consumer ID。 |
ownerId | 无 | Topic Owner ID,高级配置。 |
ownerEpoch | 无 | Topic Owner Epoch,高级配置。 |
heartbeatIntervalMs | 30000,最小 1000 | 心跳间隔,单位为毫秒。 |
endpointsSyncIntervalMs | 120000,最小 5000 | 集群端点同步间隔,单位为毫秒。 |
fileSaveDir | <user.dir>/iotdb-subscription | TsFile 临时保存目录。 |
fileSaveFsync | false | 保存 TsFile 时是否执行 fsync。 |
thriftMaxFrameSize | 由 SDK 决定 | Thrift 最大帧大小。 |
connectionTimeoutInMs | 0 | 连接超时,0 表示使用 SDK 默认行为。 |
maxPollParallelism | 1 | 最大并行 Poll 数。 |
3.1.2 Pull Consumer 特殊配置
| Builder 参数 | 默认值 | 说明 |
|---|---|---|
autoCommit | true | 是否自动提交消费进度;为 false 时必须手动提交。 |
autoCommitIntervalMs | 5000,最小 500 | 自动提交间隔,仅在 autoCommit=true 时生效。 |
3.1.3 Push Consumer 特殊配置
| Builder 参数 | 默认值 | 说明 |
|---|---|---|
ackStrategy | AckStrategy.AFTER_CONSUME | 支持 BEFORE_CONSUME 和 AFTER_CONSUME。 |
consumeListener | 始终返回 SUCCESS | 消费回调,业务代码通常需要显式提供。 |
autoPollIntervalMs | 100,最小 1 | 自动拉取间隔,单位为毫秒。 |
autoPollTimeoutMs | 10000,最小 1000 | 每次拉取的超时时间,单位为毫秒。 |
3.2 函数列表
3.2.1 Topic 管理 Session
ISubscriptionTreeSession 提供以下能力:
| 方法 | 说明 | 返回值 |
|---|---|---|
open() | 打开 Session。 | void |
createTopic(String topicName) | 使用默认配置创建 Topic。 | void |
createTopic(String topicName, Properties config) | 使用指定配置创建 Topic。 | void |
createTopicIfNotExists(...) | Topic 不存在时创建。 | void |
alterTopic(String, Properties) | 修改 Topic 配置。 | void |
alterTopicOwner(...) | 修改 Topic Owner。 | void |
dropTopic(String topicName) | 删除 Topic。 | void |
dropTopicIfExists(String topicName) | Topic 存在时删除。 | void |
getTopics() | 获取全部 Topic。 | Set<Topic> |
getTopic(String topicName) | 获取指定 Topic,不存在时为空。 | Optional<Topic> |
getSubscriptions() | 获取全部订阅关系。 | Set<Subscription> |
getSubscriptions(String topicName) | 获取指定 Topic 的订阅关系。 | Set<Subscription> |
dropSubscription(String) | 删除指定订阅关系。 | void |
dropSubscriptionIfExists(String) | 订阅关系存在时删除。 | void |
close() | 关闭 Session。 | void |
3.2.2 ISubscriptionTreePullConsumer
| 方法 | 说明 |
|---|---|
open() / close() | 打开或关闭 Consumer。 |
subscribe(String/String.../Set<String>) | 订阅一个或多个 Topic。 |
unsubscribe(String/String.../Set<String>) | 取消订阅一个或多个 Topic。 |
poll(Duration/long) | 从当前订阅的 Topic 拉取消息。 |
poll(Set<String>, Duration/long) | 从指定 Topic 集合拉取消息。 |
drainBufferedMessages() | 取出客户端已经缓冲的消息。 |
commitSync(...) | 同步提交一条或多条消息。 |
commitAsync(...) | 异步提交一条或多条消息,可指定回调。 |
seekToBeginning(String) | 将 Topic 消费位置移动到开头。 |
seekToEnd(String) | 将 Topic 消费位置移动到末尾。 |
positions(String) | 获取当前消费位置。 |
committedPositions(String) | 获取已经提交的消费位置。 |
seek(String, TopicProgress) | 将消费位置移动到指定进度。 |
seekAfter(String, TopicProgress) | 将消费位置移动到指定进度之后。 |
getConsumerId() | 获取 Consumer ID。 |
getConsumerGroupId() | 获取 Consumer Group ID。 |
allTopicMessagesHaveBeenConsumed() | 判断当前 Topic 消息是否已消费完。 |
3.2.3 ISubscriptionTreePushConsumer
| 方法或 Builder 配置 | 说明 |
|---|---|
open() / close() | 打开或关闭 Consumer。 |
subscribe(...) / unsubscribe(...) | 订阅或取消订阅一个或多个 Topic。 |
ackStrategy(AckStrategy) | 配置消息确认策略。 |
consumeListener(ConsumeListener) | 配置消费回调。 |
autoPollIntervalMs(long) | 配置自动轮询间隔。 |
autoPollTimeoutMs(long) | 配置单次拉取超时时间。 |
3.2.4 SubscriptionMessage
SubscriptionMessage 是 Consumer 获取到的基本消息单元。
| 方法 | 说明 | 返回值 |
|---|---|---|
getMessageType() | 获取消息类型。 | short |
getResultSets() | 获取 Record ResultSet。 | List<ResultSet> |
getRecordTabletIterator() | 获取 Record Tablet 迭代器。 | Iterator<Tablet> |
getTsFile() | 获取 TsFile Handler。 | SubscriptionTsFileHandler |
getCommitContext() | 获取消息提交上下文。 | SubscriptionCommitContext |
3.2.5 SubscriptionTsFileHandler
| 方法 | 说明 | 返回值 |
|---|---|---|
getFile() | 获取本地临时文件。 | File |
getPath() | 获取本地临时文件路径。 | Path |
openTreeReader() | 打开树模型 TsFile Reader。 | ITsFileTreeReader |
openTableReader() | 打开表模型 TsFile Reader。 | ITsFileReader |
copyFile(String/Path) | 将订阅文件复制到目标路径。 | Path |
moveFile(String/Path) | 将订阅文件移动到目标路径。 | Path |
deleteFile() | 删除本地订阅文件。 | Path |
4. 连接超时与资源清理
- 服务端通过心跳检测 Consumer 是否仍然活跃。
- Consumer 长时间不活跃时,服务端可以自动断开连接。
- 服务端断开连接时会同步触发该 Consumer 的取消订阅逻辑,清理订阅关系并释放资源。
- 应使用 try-with-resources 或在
finally中调用close(),确保 Session 和 Consumer 正常退出。