Go goroutine场景下如何保证Subscribe先于Publish执行?
解决方案:确保goroutine中Subscribe先于Publish执行
原代码的问题在于goroutine的调度顺序不受控,主goroutine可能在订阅goroutine还没完成Subscribe调用、还没注册好监听器的时候,就已经执行了Publish,导致消息丢失。
方案1:使用信号通道同步
新增无缓冲的信号通道,等订阅逻辑执行完成后通知主goroutine再执行发布:
useCase := New(tt.fields.storage) tt.fields.wg.Add(1) // 新增无缓冲信号通道,用于通知订阅完成 subDone := make(chan struct{}) go func() { defer tt.fields.wg.Done() ch, _, err := useCase.Subscribe(tt.args.ctx, tt.args.message.TopicName) // 订阅调用完成,发送信号通知主goroutine可以发布 close(subDone) require.NoError(t, err) message, ok := <-ch if !ok { close(ch) return } assert.Equal(t, tt.want, message) }() // 阻塞等待订阅完成信号,再执行发布 <-subDone err := useCase.Publish(tt.args.ctx, tt.args.message) tt.fields.wg.Wait()
方案2:使用单独的WaitGroup同步订阅逻辑
也可以用专门的WaitGroup管控订阅流程的同步:
useCase := New(tt.fields.storage) tt.fields.wg.Add(1) var subWg sync.WaitGroup subWg.Add(1) go func() { defer tt.fields.wg.Done() ch, _, err := useCase.Subscribe(tt.args.ctx, tt.args.message.TopicName) // 订阅完成,标记subWg完成 subWg.Done() require.NoError(t, err) message, ok := <-ch if !ok { close(ch) return } assert.Equal(t, tt.want, message) }() // 等待订阅完成 subWg.Wait() err := useCase.Publish(tt.args.ctx, tt.args.message) tt.fields.wg.Wait()
注意事项
- 如果
Subscribe方法内部是异步注册监听器、调用时不会等注册完成就返回,你需要把信号发送/subWg.Done()的位置挪到Subscribe内部监听器注册完成的逻辑节点,确保真正注册成功之后再通知发布。 - 无缓冲通道的
close操作是广播通知,即便后续有多个监听方也可以正常收到信号,比发送单个struct{}适用场景更广。
内容的提问来源于stack exchange,提问作者Ming
相关产品推荐
相关产品推荐

