通道阻塞时如何解锁?Goroutine锁与消息顺序保障方案咨询
sendTo 是可在不同位置作为 goroutine 调用的函数,最初的实现如下:
func (s *Sender) sendTo() { for { s.Lock() if s.next == len(s.msgs) { s.Unlock() break } // prepare message... msg := s.msgs[s.next] s.next++ s.ch <- msg s.Unlock() } }
由于s的成员会被不同 goroutine 访问(仅 sendTo 中使用s.ch),因此用锁保护这些成员。但当s.ch已满时,sendTo 会长期持有锁,阻塞其他 goroutine 访问s的成员。
为解决锁持有时间过长的问题,修改后的代码如下:
func (s *Sender) sendTo() { for { s.Lock() if s.next == len(s.msgs) { s.Unlock() break } // prepare message... msg := s.msgs[s.next] s.next++ s.Unlock() s.ch <- msg } }
但修改后的代码无法保证发送到s.ch的消息顺序与s.msgs一致——多个goroutine并发取出消息后,通道的发送顺序取决于goroutine调度,会打乱原有的消息顺序。现咨询是否存在类似条件变量的方式,可先阻塞通道再解锁,同时保证消息顺序。
可行解决方案
方案1:单发送goroutine + 任务队列(推荐)
最简洁的方式是拆分职责:让一个单独的goroutine负责按顺序将消息发送到s.ch,其他sendTo goroutine只需要安全地取出消息,放到一个任务通道中。这样既避免了锁长期持有,又天然保证了发送顺序。
首先给Sender结构体添加一个任务通道:
type Sender struct { sync.Mutex msgs []Message next int ch chan Message taskCh chan Message // 新增任务通道 } // 初始化时启动单发送goroutine func NewSender(msgs []Message, chSize int) *Sender { s := &Sender{ msgs: msgs, ch: make(chan Message, chSize), taskCh: make(chan Message, len(msgs)), } // 启动专门的发送goroutine go func() { for msg := range s.taskCh { s.ch <- msg } close(s.ch) // 所有消息发送完成后关闭通道 }() return s }
然后修改sendTo函数:
func (s *Sender) sendTo() { for { s.Lock() if s.next == len(s.msgs) { s.Unlock() break } msg := s.msgs[s.next] s.next++ s.Unlock() // 将消息放到任务通道,由单goroutine负责发送 s.taskCh <- msg } // 所有sendTo完成后关闭任务通道(需额外同步,比如用sync.WaitGroup) }
这种方式下,锁只用来保护s.next和s.msgs的访问,持有时间极短;单goroutine消费任务通道,保证消息按放入顺序发送到s.ch,完美解决两个问题。
方案2:用条件变量同步发送顺序
如果不想新增goroutine,可以用条件变量来确保每个goroutine按取出消息的索引顺序发送:
首先给Sender结构体添加条件变量和当前发送索引:
type Sender struct { sync.Mutex cond *sync.Cond msgs []Message next int currentSendIdx int // 当前应该发送的消息索引 ch chan Message } // 初始化条件变量 func NewSender(msgs []Message, chSize int) *Sender { s := &Sender{ msgs: msgs, ch: make(chan Message, chSize), currentSendIdx: 0, } s.cond = sync.NewCond(&s.Mutex) return s }
修改sendTo函数:
func (s *Sender) sendTo() { for { s.Lock() if s.next == len(s.msgs) { s.Unlock() break } // 获取当前要发送的消息索引和内容 idx := s.next msg := s.msgs[idx] s.next++ s.Unlock() // 等待直到轮到自己发送 s.Lock() for idx != s.currentSendIdx { s.cond.Wait() // Wait会自动释放锁,被唤醒时重新获取锁 } // 发送消息后更新索引并通知其他goroutine s.ch <- msg s.currentSendIdx++ s.cond.Broadcast() // 通知所有等待的goroutine检查是否轮到自己 s.Unlock() } }
这种方式下,每个goroutine取出消息后,会等待自己的索引等于当前应该发送的索引,发送完成后更新索引并通知其他goroutine。由于条件变量的Wait会自动释放锁,所以不会在等待时持有锁;发送消息时虽然持有锁,但只有当通道有空间时才会执行发送,阻塞时间可控,同时保证了发送顺序。
内容的提问来源于stack exchange,提问作者ihnfsa

