如何列出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
相关产品推荐
相关产品推荐

