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

如何列出NATS JetStream消息并识别未被确认的消息?

列出NATS Stream中未确认消息的可行方案
  • 使用NATS CLI工具
    这是最直接的方式,官方CLI提供了专门查看未确认消息的命令:

    • 查看消费者未确认消息统计:nats consumer info <stream-name> <consumer-name>,输出中会包含未确认消息数量及单条消息的序列号、过期时间。
    • 直接列出所有未确认消息详情:
      nats jetstream ack pending <stream-name> <consumer-name>
      
  • Go语言(官方新版jetstream库)
    旧版jsm.go已被官方弃用,推荐使用github.com/nats-io/nats.go/jetstream库,通过FetchPending方法直接获取未确认消息:

    package main
    
    import (
        "context"
        "fmt"
        "log"
        "github.com/nats-io/nats.go"
        "github.com/nats-io/nats.go/jetstream"
    )
    
    func main() {
        nc, err := nats.Connect("nats://localhost:4222")
        if err != nil {
            log.Fatal(err)
        }
        defer nc.Close()
    
        ctx := context.Background()
        js, err := jetstream.New(nc)
        if err != nil {
            log.Fatal(err)
        }
    
        cons, err := js.Consumer(ctx, "my-stream", "my-consumer")
        if err != nil {
            log.Fatal(err)
        }
    
        // 批量获取未确认消息,第二个参数为单次拉取数量
        pendingMsgs, err := cons.FetchPending(ctx, 50)
        if err != nil {
            log.Fatal(err)
        }
    
        for msg := range pendingMsgs.Messages() {
            fmt.Printf("未确认消息流序列号: %d, 内容: %s\n", msg.Sequence.Stream, string(msg.Data()))
            // 注意:仅查看时不要调用msg.Ack(),避免误确认
        }
    }
    
  • Python(nats.py库)
    通过nats.py的JetStream API可获取未确认消息统计,也能直接拉取未确认消息:

    import asyncio
    import nats
    from nats.js import JetStream
    
    async def main():
        nc = await nats.connect("nats://localhost:4222")
        js = JetStream(nc)
    
        # 获取消费者未确认消息统计
        consumer_info = await js.consumer_info(stream="my-stream", consumer="my-consumer")
        print(f"未确认消息总数: {consumer_info.num_pending}")
    
        # 拉取未确认消息
        sub = await js.pull_subscribe(
            subject=">", stream="my-stream", consumer="my-consumer", batch=20
        )
        async for msg in sub.messages:
            print(f"未确认消息流序列号: {msg.stream_seq}, 内容: {msg.data.decode()}")
            # 仅查看时不要调用msg.ack()
    
        await nc.close()
    
    if __name__ == "__main__":
        asyncio.run(main())
    
  • Admin API方式
    官方Admin API其实支持该操作,通过GET /js/api/consumers/<stream-name>/<consumer-name>接口,返回的JSON响应中:

    • num_pending字段为未确认消息总数
    • pending数组包含每条未确认消息的序列号(seq)和过期时间(expires)

内容的提问来源于stack exchange,提问作者Jonathan ch

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 03:10:56