Go语言NATS客户端封装:自定义Handler的订阅方法阻塞实现问题
问题解决:Go封装NATS客户端订阅阻塞与消息接收问题
核心问题分析
- 订阅被提前取消:原
Subscribe方法执行完毕后,defer subscription.Unsubscribe()会立即取消订阅,导致后续消息无法被接收。 - 订阅未完成就发布:测试中
Publish在Subscribe后立即执行,此时NATS客户端可能还未将订阅请求同步到服务器,消息直接丢失。 - 阻塞逻辑错误:原代码的WaitGroup仅等待单条消息,且未解决订阅提前取消的根本问题。
解决方案
1. 修正Subscribe方法:保持阻塞并避免提前取消订阅
根据需求,提供两种阻塞实现:
方式一:永久阻塞(持续接收消息)
让方法一直阻塞,直到程序主动退出,适合长期运行的服务:
func (busConfig BusConfig) Subscribe(subject string, handler func(msg []byte)) error { fmt.Println("Subscribing on : ", subject) // 强制同步订阅请求到服务器,确保订阅生效 if err := busConfig.conn.Flush(); err != nil { return fmt.Errorf("flush subscribe request failed: %w", err) } subscription, err := busConfig.conn.Subscribe(subject, func(msg *nats.Msg) { // 若handler耗时较长,可保留goroutine包裹;简单处理则直接调用 go handler(msg.Data) }) if err != nil { return fmt.Errorf("create subscription failed: %w", err) } defer subscription.Unsubscribe() // 永久阻塞,避免方法退出触发取消订阅 select{} }
方式二:阻塞直到收到第一条消息
适合仅需接收单条消息的场景:
func (busConfig BusConfig) Subscribe(subject string, handler func(msg []byte)) error { fmt.Println("Subscribing on : ", subject) if err := busConfig.conn.Flush(); err != nil { return fmt.Errorf("flush subscribe request failed: %w", err) } done := make(chan struct{}) subscription, err := busConfig.conn.Subscribe(subject, func(msg *nats.Msg) { handler(msg.Data) close(done) // 收到消息后触发退出 }) if err != nil { return fmt.Errorf("create subscription failed: %w", err) } defer subscription.Unsubscribe() // 阻塞直到收到第一条消息 <-done return nil }
2. 修正测试用例:确保订阅生效后再发布
func TestLifeCycleEvent(t *testing.T) { busClient := GetBusClient() // 用goroutine执行订阅(若Subscribe是永久阻塞) go func() { if err := busClient.Subscribe(SUBJECT, func(input []byte) { fmt.Println("Life cycle event received :", string(input)) }); err != nil { t.Fatalf("subscribe error: %v", err) } }() // 等待订阅请求同步到服务器 if err := busClient.conn.Flush(); err != nil { t.Fatalf("flush connection failed: %v", err) } // 发布消息 if err := busClient.Publish(SUBJECT, []byte("complete notification")); err != nil { t.Fatalf("publish error: %v", err) } // 短暂等待消息处理完成(测试场景用,生产环境可通过信号量替代) time.Sleep(100 * time.Millisecond) }
关键注意事项
- 禁止提前取消订阅:必须让
Subscribe方法保持阻塞,避免defer触发订阅取消。 - 强制同步订阅请求:
conn.Flush()是确保订阅在发布前生效的关键,避免消息丢失。 - Handler goroutine按需使用:仅当Handler执行耗时操作时才用goroutine包裹,否则直接调用即可。
内容的提问来源于stack exchange,提问作者Ravat Tailor
相关产品推荐
相关产品推荐

