如何在Java客户端实现无Durable的JetStream批量拉取消息?
问题解答
核心结论:无需Durable订阅可实现需求
可以不使用PullSubscribeOptions.durable()实现拉取全部消息或最后一批消息的功能,但必须手动指定消费起始位置——因为非Durable订阅是临时消费者,服务器不会保存其消费进度,默认每次都从流的起始位置发送消息,这就是你遇到重复批次的原因。通过明确起始位置,就能避免重复拉取,同时实现你的两个需求。
分场景实现方案
场景1:拉取流中全部消息(支持重复执行)
每次拉取时指定从流的第一条消息开始,分批拉取直到无新消息:
public static void main(String[] args) throws IOException, InterruptedException, JetStreamApiException, TimeoutException { String userName = "local"; String password = "password"; String host = "0.0.0.0"; String port = "37733"; String streamName = "test_stream"; String subjectName = "test_subject_1"; int BATCH_SIZE = 100; try (Connection nc = Nats.connect("nats://" + userName + ":" + password + "@" + host + ":" + port)) { JetStream js = nc.jetStream(); while (true) { // 创建临时订阅,指定从流的起始位置开始拉取 PullSubscribeOptions pullOptions = PullSubscribeOptions.builder() .stream(streamName) .startAt(StartPosition.FIRST) .build(); try (JetStreamSubscription sub = js.subscribe(subjectName, pullOptions)) { boolean hasMore; do { List<Message> messages = sub.fetch(BATCH_SIZE, Duration.ofMillis(1000)); hasMore = !messages.isEmpty(); for (Message m : messages) { System.out.println(m); m.ack(); } System.out.println("本次拉取消息数:" + messages.size()); } while (hasMore); } System.out.println("全量消息拉取完成,等待下一轮执行..."); Thread.sleep(5000); // 自定义重复执行间隔 } } }
场景2:拉取最后一批消息(指定批量大小)
先获取流的最新序列,计算起始位置后拉取指定数量的最新消息:
public static void main(String[] args) throws IOException, InterruptedException, JetStreamApiException, TimeoutException { String userName = "local"; String password = "password"; String host = "0.0.0.0"; String port = "37733"; String streamName = "test_stream"; String subjectName = "test_subject_1"; int BATCH_SIZE = 100; try (Connection nc = Nats.connect("nats://" + userName + ":" + password + "@" + host + ":" + port)) { JetStream js = nc.jetStream(); while (true) { // 获取流的最新状态,计算起始序列 StreamInfo streamInfo = js.streamInfo(streamName); long lastSeq = streamInfo.getState().getLastSequence(); long startSeq = Math.max(1, lastSeq - BATCH_SIZE + 1); // 避免序列小于1 // 创建临时订阅,从计算出的起始位置拉取 PullSubscribeOptions pullOptions = PullSubscribeOptions.builder() .stream(streamName) .startAtSequence(startSeq) .build(); try (JetStreamSubscription sub = js.subscribe(subjectName, pullOptions)) { List<Message> messages = sub.fetch(BATCH_SIZE, Duration.ofMillis(1000)); for (Message m : messages) { System.out.println(m); m.ack(); } System.out.println("本次拉取最新消息数:" + messages.size()); } System.out.println("等待下一次拉取最新消息..."); Thread.sleep(5000); } } }
关键注意事项
- 避免循环内重复创建订阅:你原代码在while循环里每次都新建订阅,这会产生大量临时消费者,效率极低。优化后的代码将订阅创建放在循环内但用try-with-resources自动关闭,或者如果是长期重复拉取,也可以考虑复用订阅并重置起始位置。
- 非Durable订阅的本质:临时消费者不会在服务器留存任何状态,所以每次拉取必须明确起始位置,否则默认从流开头返回消息,导致重复拉取。
- nats-cli的实现逻辑:nats-cli的
jetstream pull命令不带durable时,也是通过--first、--last等参数指定起始位置,和上述代码逻辑一致。
内容的提问来源于stack exchange,提问作者Janith Priyankara
相关产品推荐
相关产品推荐

