Golang基于Channel与Select实现Pub/Sub并发的异常问题求助
问题:Golang并发Pub/Sub模型中已移除的订阅者仍输出空消息
我用Golang实现了一个基于并发的Pub/Sub模型,代码有时能正常执行,但偶尔会出现异常输出:
异常输出
message received on channel 2: Hello World message received on channel 3: Hello World message received on channel 1: Hello World message received on channel 1: subscriber 1's context cancelled message received on channel 3: Only channels 2 and 3 should print this message received on channel 2: Only channels 2 and 3 should print this
关键问题是:订阅者1在调用RemoveSubscriber移除后,仍打印了message received on channel 1: 。订阅者1应该只接收第一条消息,之后其context被取消、goroutine退出。目前推测是订阅者的dctx被取消前,从subscriber.out通道收到了消息(通道关闭后接收会得到零值)。
预期执行结果:
通道1仅打印一次消息,之后context被取消,goroutine退出,不再接收任何消息。
main.go
package main import ( "fmt" "net/http" ) func main() { var err error publisher := NewPublisher() publisher.AddSubscriber() publisher.AddSubscriber() publisher.AddSubscriber() publisher.Start() err = publisher.Publish("Hello World") if err != nil { fmt.Printf("could not publish: %v\n", err) } err = publisher.RemoveSubscriber(1) if err != nil { fmt.Printf("could not remove subscriber: %v\n", err) } err = publisher.Publish("Only channels 2 and 3 should print this") if err != nil { fmt.Printf("could not publish: %v\n", err) } // 保持服务运行 http.ListenAndServe(":8080", nil) }
publisher.go
package main import ( "context" "errors" "fmt" "sync" ) // Publisher 将in通道收到的消息发送给所有订阅者 type Publisher struct { sequence uint in chan string subscribers map[uint]*Subscriber sync.RWMutex ctx context.Context cancel *context.CancelFunc } // NewPublisher 返回一个空订阅者的Publisher实例 // 必须调用Start()方法后才能开始发布消息 func NewPublisher() *Publisher { ctx, cancel := context.WithCancel(context.Background()) return &Publisher{subscribers: map[uint]*Subscriber{}, ctx: ctx, cancel: &cancel} } // AddSubscriber 创建新订阅者并开始监听Publisher的消息 func (p *Publisher) AddSubscriber() { dctx, cancel := context.WithCancel(p.ctx) p.Lock() nextId := p.sequence + 1 subscriber := NewSubscriber(nextId, dctx, &cancel) p.subscribers[nextId] = subscriber p.sequence = p.sequence + 1 p.Unlock() go func() { for { select { case <-p.ctx.Done(): fmt.Printf("parent context cancelled\n") (*subscriber.cancel)() return case <-dctx.Done(): fmt.Printf("subscriber %d's context cancelled\n", subscriber.id) return case msg := <-subscriber.out: fmt.Printf("message received on channel %d: %s\n", subscriber.id, msg) } } }() } // Publish 将消息发送给所有订阅者 func (p *Publisher) Publish(msg string) error { // 未启动则返回错误 if p.in == nil { return errors.New("publisher not started yet") } // 无订阅者则返回错误 p.RLock() if len(p.subscribers) == 0 { return errors.New("no subscribers to receive the message") } p.RUnlock() // 发送消息到in通道,加锁防止并发修改 p.Lock() p.in <- msg p.Unlock() return nil } // Start 初始化in通道,使Publisher可以接收并发布消息 func (p *Publisher) Start() { in := make(chan string) p.in = in go func() { for { select { case <-p.ctx.Done(): fmt.Printf("done called on publisher\n") return case msg := <-p.in: p.RLock() for _, subscriber := range p.subscribers { subscriber.out <- msg } p.RUnlock() } } }() } // Stop 终止Publisher的消息监听 func (p *Publisher) Stop() { (*p.cancel)() } // RemoveSubscriber 根据ID移除订阅者,找不到则返回错误 func (p *Publisher) RemoveSubscriber(id uint) error { p.Lock() defer p.Unlock() subscriber, ok := p.subscribers[id] if !ok { return errors.New("could not find subscriber") } (*subscriber.cancel)() delete(p.subscribers, id) close(subscriber.out) return nil }
问题分析与修复方案
问题根源
出现空消息的核心原因是:在RemoveSubscriber中,你先取消订阅者的context,紧接着关闭了subscriber.out通道。但订阅者goroutine的select分支中,context取消信号和通道读取是平等竞争的,有可能在context取消逻辑执行前,goroutine先读取到了关闭通道返回的零值(空字符串),从而打印出异常内容。
修复步骤
- 删除手动关闭通道的操作:不需要主动关闭
subscriber.out,订阅者goroutine会在context取消后自动退出,通道会被Go的垃圾回收机制处理,避免了读取零值的问题。 - 增强订阅者goroutine的可靠性:在读取消息后,额外检查context状态,确保即使消息在context取消前到达,也不会被处理。
修改后的关键代码
修改RemoveSubscriber方法
移除close(subscriber.out)语句:
// RemoveSubscriber 根据ID移除订阅者,找不到则返回错误 func (p *Publisher) RemoveSubscriber(id uint) error { p.Lock() defer p.Unlock() subscriber, ok := p.subscribers[id] if !ok { return errors.New("could not find subscriber") } (*subscriber.cancel)() delete(p.subscribers, id) // 移除关闭通道的操作,避免读取零值 // close(subscriber.out) return nil }
优化订阅者goroutine逻辑
在处理消息前再次检查context状态,确保不会处理取消后的消息:
go func() { for { select { case <-p.ctx.Done(): fmt.Printf("parent context cancelled\n") (*subscriber.cancel)() return case <-dctx.Done(): fmt.Printf("subscriber %d's context cancelled\n", subscriber.id) return case msg := <-subscriber.out: // 再次检查context是否已取消,避免处理延迟到达的消息 select { case <-dctx.Done(): return default: fmt.Printf("message received on channel %d: %s\n", subscriber.id, msg) } } } }()
修复原理
- 不再关闭订阅者通道,彻底避免了读取通道关闭后零值的情况。
- 订阅者goroutine会在context取消信号触发后立即退出,不会再处理任何后续消息。
- 额外的context检查确保即使消息在取消信号前已经到达通道,也会直接退出,不执行打印逻辑。
内容的提问来源于stack exchange,提问作者Zaid Sheikh
相关产品推荐
相关产品推荐

