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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 12:57:04