观察者模式在Go语言里实现起来并不复杂,核心思路就是定义一个统一的观察者接口,让被观察者把观察者实例集中管理,在自身状态变化时逐个通知。这个模式最典型的应用场景包括消息推送、配置变更、缓存失效、订单状态流转等。实现时看似只要遍历回调就够了,但一旦加入并发、慢任务和动态注销,细节就会成倍增加。本文从同步版本开始,逐步把通知链路改造成带缓冲队列的异步分发,并讨论锁策略、panic隔离和观察者注销这些工程化问题。

一、同步观察者:用接口建立最小通知链路
实现观察者模式的第一步,是让所有观察者遵循同一个方法约定。在Go里最直接的做法就是定义一个Observer接口,只暴露一个Update方法。被观察者内部不关心具体实现,只保存一组Observer接口实例。这样做的好处是,新增邮件通知、短信通知还是日志记录,都不需要修改Subject的代码,只要实现Update方法并注册进去即可。
下面是一个最简单的同步版本。被观察者用切片保存观察者,通知时按注册顺序逐个调用Update。这个版本适合原型验证或者观察者执行速度都很快的场景。
package main
import "fmt"
type Observer interface {
Update(event string)
}
type Subject struct {
observers []Observer
}
func (s *Subject) AddObserver(o Observer) {
s.observers = append(s.observers, o)
}
func (s *Subject) RemoveObserver(o Observer) {
for i, ob := range s.observers {
if ob == o {
s.observers = append(s.observers[:i], s.observers[i+1:]...)
break
}
}
}
func (s *Subject) Notify(event string) {
for _, ob := range s.observers {
ob.Update(event)
}
}
type EmailObserver struct {
name string
}
func (e *EmailObserver) Update(event string) {
fmt.Printf("[%s] email observer received event: %s\n", e.name, event)
}
func main() {
subject := new(Subject)
email := new(EmailObserver)
email.name = "email"
subject.AddObserver(email)
sms := new(EmailObserver)
sms.name = "sms"
subject.AddObserver(sms)
subject.Notify("order created")
}
同步实现的优点在于代码直观,调用栈清晰,排错容易。但它的缺点也很明显:Notify会按照顺序逐个执行观察者,如果第一个观察者调用外部接口耗时500毫秒,那么后面的所有观察者都要等这500毫秒。比如订单创建接口里同步发送邮件和短信,一旦短信网关响应变慢,订单接口的响应时间就会被拉长。因此,当观察者中出现了耗时的外部依赖,就应当把通知过程异步化。
二、异步观察者:用通道和缓冲队列解除阻塞
要让业务主流程不再等待观察者执行完毕,常见做法是给Subject增加一个带缓冲的事件通道。业务代码调用Notify时,只是把事件投递到通道里,如果缓冲未满就立即返回。真正的通知动作由一个独立的goroutine负责,从通道里读取事件,再分发给所有观察者。这样的设计把状态变更和通知执行拆成了两个阶段。
先定义一个事件结构,再改造Subject。事件结构体里除了类型字段,还可以带上一个interface{}类型的载荷,方便后续扩展不同业务数据。分发循环在Start方法中启动,通过select同时监听事件通道和停止通道。
package main
import (
"fmt"
"sync"
)
type Event struct {
Type string
Payload interface{}
}
type Observer interface {
Update(event Event)
}
type Subject struct {
mu sync.Mutex
observers map[string]Observer
eventCh chan Event
stopCh chan struct{}
stopOnce sync.Once
}
func NewSubject() *Subject {
s := new(Subject)
s.observers = make(map[string]Observer)
s.eventCh = make(chan Event, 64)
s.stopCh = make(chan struct{})
return s
}
func (s *Subject) Start() {
go func() {
for {
select {
case event := <-s.eventCh:
s.dispatch(event)
case <-s.stopCh:
return
}
}
}()
}
func (s *Subject) dispatch(event Event) {
s.mu.Lock()
snapshot := make([]Observer, 0, len(s.observers))
for _, ob := range s.observers {
snapshot = append(snapshot, ob)
}
s.mu.Unlock()
for _, ob := range snapshot {
ob.Update(event)
}
}
func (s *Subject) Notify(event Event) {
select {
case s.eventCh <- event:
default:
fmt.Println("event channel is full, drop event:", event.Type)
}
}
func (s *Subject) AddObserver(id string, o Observer) {
s.mu.Lock()
defer s.mu.Unlock()
s.observers[id] = o
}
func (s *Subject) RemoveObserver(id string) {
s.mu.Lock()
defer s.mu.Unlock()
delete(s.observers, id)
}
func (s *Subject) Stop() {
s.stopOnce.Do(func() {
close(s.stopCh)
})
}
这里选择共享一个分发goroutine,而不是为每个观察者单独启动goroutine,主要是为了控制复杂度。共享循环的好处是注销观察者非常安全,直接删除map中的条目即可,不需要处理某个观察者专属通道的关闭和发送竞态。它的局限在于,dispatch内部仍然是串行调用Update,一个特别慢的观察者会阻塞后续事件的分发。如果业务上要求观察者之间完全隔离,可以再升级成每个观察者一个通道和一个goroutine的模型,但那时要额外处理goroutine退出、通道关闭顺序以及消息积压的问题。
三、并发安全与工程化细节
在并发场景下,最重要的一点是不要长期持锁调用外部代码。上面的dispatch方法先加锁拷贝观察者列表,随后立即解锁,再遍历快照执行通知。这样做有两个好处:一是避免在持有sync.Mutex期间调用Update,防止外部代码回调AddObserver或RemoveObserver造成死锁;二是缩短临界区,减少对注册、注销操作的阻塞。对于读多写少的场景,也可以把sync.Mutex换成sync.RWMutex,通知时使用读锁,注册和注销时使用写锁。但要注意,通知读锁下仍然不建议直接调用外部代码,最好还是快照后解锁。
另一个容易忽略的问题是观察者自身panic导致的连锁影响。分发goroutine里如果某个Update抛出panic,整个通知循环会直接崩溃,后续事件全部丢失。更稳妥的做法是在调用每个观察者时增加一层恢复逻辑,把单个观察者的异常隔离住。下面的代码可以放在Subject中作为统一的通知出口。
func (s *Subject) safeNotify(ob Observer, event Event) {
defer func() {
if r := recover(); r != nil {
fmt.Printf("observer panic recovered: %v\n", r)
}
}()
ob.Update(event)
}
然后dispatch里把ob.Update(event)替换成s.safeNotify(ob, event)即可。除了panic问题,还要考虑注销观察者的时机。当前共享分发模型的注销操作非常轻量,只要加锁删除map条目,后续分发就不再包含该观察者。如果未来改成每个观察者一个goroutine,注销时必须给worker发送停止信号,而且不能从发送侧贸然关闭通道,否则很容易出现向已关闭通道发送数据的panic。因此在设计上,优先用共享分发循环,等确实遇到单个慢观察者拖慢整体的情况,再针对性地引入独立队列。
四、完整示例与使用建议
把上面的片段整合起来,可以形成一个直接运行的示例。下面代码演示了两个具体观察者注册、事件投递和停止分发的完整过程。Event.Payload使用interface{}是为了保持示例简洁,实际项目中更推荐替换成具体的业务事件结构体,或者利用Go泛型来获得编译期类型安全。
package main
import (
"fmt"
"sync"
)
type Event struct {
Type string
Payload interface{}
}
type Observer interface {
Update(event Event)
}
type Subject struct {
mu sync.Mutex
observers map[string]Observer
eventCh chan Event
stopCh chan struct{}
stopOnce sync.Once
}
func NewSubject() *Subject {
s := new(Subject)
s.observers = make(map[string]Observer)
s.eventCh = make(chan Event, 64)
s.stopCh = make(chan struct{})
return s
}
func (s *Subject) Start() {
go func() {
for {
select {
case event := <-s.eventCh:
s.dispatch(event)
case <-s.stopCh:
return
}
}
}()
}
func (s *Subject) dispatch(event Event) {
s.mu.Lock()
snapshot := make([]Observer, 0, len(s.observers))
for _, ob := range s.observers {
snapshot = append(snapshot, ob)
}
s.mu.Unlock()
for _, ob := range snapshot {
s.safeNotify(ob, event)
}
}
func (s *Subject) safeNotify(ob Observer, event Event) {
defer func() {
if r := recover(); r != nil {
fmt.Printf("observer panic recovered: %v\n", r)
}
}()
ob.Update(event)
}
func (s *Subject) Notify(event Event) {
select {
case s.eventCh <- event:
default:
fmt.Println("event channel is full, drop event:", event.Type)
}
}
func (s *Subject) AddObserver(id string, o Observer) {
s.mu.Lock()
defer s.mu.Unlock()
s.observers[id] = o
}
func (s *Subject) RemoveObserver(id string) {
s.mu.Lock()
defer s.mu.Unlock()
delete(s.observers, id)
}
func (s *Subject) Stop() {
s.stopOnce.Do(func() {
close(s.stopCh)
})
}
type EmailObserver struct {
name string
}
func (e *EmailObserver) Update(event Event) {
fmt.Printf("[%s] email observer received event: %s\n", e.name, event.Type)
}
type SMSObserver struct {
name string
}
func (s *SMSObserver) Update(event Event) {
fmt.Printf("[%s] sms observer received event: %s\n", s.name, event.Type)
}
func main() {
subject := NewSubject()
email := new(EmailObserver)
email.name = "email-observer"
subject.AddObserver("email", email)
sms := new(SMSObserver)
sms.name = "sms-observer"
subject.AddObserver("sms", sms)
subject.Start()
subject.Notify(Event{Type: "user.registered", Payload: "user_id=1001"})
subject.Notify(Event{Type: "order.created", Payload: "order_id=202"})
subject.Stop()
}
如果在实际业务中事件吞吐量较高,可以适当调大eventCh的缓冲容量,并结合监控观察队列是否频繁写满。一旦写满,Notify当前采用丢弃策略,这种策略适合允许丢失通知但对时延敏感的系统。如果业务不允许丢消息,则不能丢弃,而是应该阻塞等待或者落盘后异步重试。另一个实践建议是,避免在Update方法内部执行长时间阻塞操作,观察者自身仍然需要对外部调用设置超时控制。这样观察者模式才能保持轻量,不会随着观察者数量增加而变成新的性能瓶颈。
总体来看,Go实现观察者模式并不需要依赖复杂框架。先用接口和切片建立最小可用版本,再用通道把同步通知改造成异步分发,最后补上锁策略、panic恢复和注销管理,就能得到一个结构清晰、扩展方便的观察者模块。你可以基于这套骨架继续加入重试机制、事件优先级、消费确认等能力,但它已经足以支撑大多数中小型Go项目的事件通知需求。