You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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);
        }
    }
}

关键注意事项

  1. 避免循环内重复创建订阅:你原代码在while循环里每次都新建订阅,这会产生大量临时消费者,效率极低。优化后的代码将订阅创建放在循环内但用try-with-resources自动关闭,或者如果是长期重复拉取,也可以考虑复用订阅并重置起始位置。
  2. 非Durable订阅的本质:临时消费者不会在服务器留存任何状态,所以每次拉取必须明确起始位置,否则默认从流开头返回消息,导致重复拉取。
  3. nats-cli的实现逻辑:nats-cli的jetstream pull命令不带durable时,也是通过--first、--last等参数指定起始位置,和上述代码逻辑一致。

内容的提问来源于stack exchange,提问作者Janith Priyankara

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 20:45:59