Appearance
观察者模式
观察者模式(Observer)是行为型模式中最常用之一,它的另一个名字是「发布-订阅(Publish-Subscribe)」。意图是:定义对象间一对多的依赖关系,当一个对象状态变化时,所有依赖它的对象都会收到通知。本章我们看 Go 用 channel 和回调两种方式实现观察者,并对比它们的适用场景。
一、观察者模式意图:发布-订阅
1. 意图
定义对象之间的一对多依赖关系,以便当一个对象的状态发生变化时,所有依赖于它的对象都会得到通知并自动更新。
经典场景:
- GUI 事件系统:按钮被点击,多个监听器被触发。
- 消息总线:业务事件发布后,日志、监控、统计等多个订阅者各自处理。
- 数据绑定:模型变化时视图自动刷新。
- 分布式事件通知:MQ 的消费模型。
2. 角色定义
- Subject(主题/被观察者):维护观察者列表,状态变化时通知所有观察者。
- Observer(观察者):定义一个更新接口,收到通知时执行自己的逻辑。
- ConcreteSubject / ConcreteObserver:具体的实现。
Java 经典写法是 Subject 持有 List<Observer>,notify() 时遍历调用 observer.update()。Go 有两种主流写法:回调和 channel。
二、Go 实现:回调函数实现观察者
最直白的实现:Subject 持有一组回调函数,状态变化时遍历调用。
go
package main
import (
"fmt"
"sync"
)
// 事件类型
type Event struct {
Topic string
Payload any
}
// 观察者就是一个回调函数
type Observer func(Event)
// Subject
type Subject struct {
mu sync.RWMutex
observers []Observer
}
func NewSubject() *Subject {
return &Subject{}
}
// Subscribe 注册观察者
func (s *Subject) Subscribe(o Observer) {
s.mu.Lock()
defer s.mu.Unlock()
s.observers = append(s.observers, o)
}
// Notify 通知所有观察者
func (s *Subject) Notify(e Event) {
s.mu.RLock()
observers := make([]Observer, len(s.observers))
copy(observers, s.observers)
s.mu.RUnlock()
for _, o := range observers {
o(e) // 注意:同步调用
}
}
func main() {
sub := NewSubject()
sub.Subscribe(func(e Event) {
fmt.Printf("[日志] %s: %v\n", e.Topic, e.Payload)
})
sub.Subscribe(func(e Event) {
fmt.Printf("[监控] %s: %v\n", e.Topic, e.Payload)
})
sub.Subscribe(func(e Event) {
fmt.Printf("[统计] %s: %v\n", e.Topic, e.Payload)
})
sub.Notify(Event{Topic: "order.created", Payload: "ORD-001"})
}几个关键点:
- 用
sync.RWMutex保护观察者列表,支持并发订阅和通知。 Notify时先copy一份观察者列表再遍历,避免回调中再次Subscribe导致死锁。- 同步调用:所有观察者按顺序执行,一个慢会拖累整个
Notify。
1. 异步通知
如果观察者可能很慢,可以让通知异步进行:
go
package main
import (
"fmt"
"sync"
"time"
)
type Event struct {
Topic string
Payload any
}
type Observer func(Event)
type AsyncSubject struct {
mu sync.RWMutex
observers []Observer
}
func (s *AsyncSubject) Subscribe(o Observer) {
s.mu.Lock()
defer s.mu.Unlock()
s.observers = append(s.observers, o)
}
func (s *AsyncSubject) Notify(e Event) {
s.mu.RLock()
defer s.mu.RUnlock()
var wg sync.WaitGroup
for _, o := range s.observers {
wg.Add(1)
go func(o Observer) {
defer wg.Done()
o(e)
}(o)
}
// 可以选择等或不等:wg.Wait()
go wg.Wait() // 这里不阻塞主流程
}
func main() {
sub := &AsyncSubject{}
sub.Subscribe(func(e Event) {
time.Sleep(100 * time.Millisecond)
fmt.Println("[A] 处理完", e.Topic)
})
sub.Subscribe(func(e Event) {
time.Sleep(50 * time.Millisecond)
fmt.Println("[B] 处理完", e.Topic)
})
start := time.Now()
sub.Notify(Event{Topic: "test"})
fmt.Println("Notify 立即返回,耗时:", time.Since(start))
time.Sleep(200 * time.Millisecond) // 等异步完成
}异步通知的代价:错误处理复杂、顺序不可控、资源消耗大。是否异步取决于业务——日志、监控这种可以异步;订单状态变更这种需要确认的,最好同步。
2. 取消订阅
实际项目里通常需要取消订阅,避免内存泄漏。可以把订阅者抽象成带 ID 的结构:
go
package main
import (
"fmt"
"sync"
)
type Event struct{ Topic string }
type Handler func(Event)
type Subscription struct {
id int64
handler Handler
}
type EventBus struct {
mu sync.RWMutex
nextID int64
subs map[string][]*Subscription
}
func NewEventBus() *EventBus {
return &EventBus{subs: map[string][]*Subscription{}}
}
func (b *EventBus) Subscribe(topic string, h Handler) *Subscription {
b.mu.Lock()
defer b.mu.Unlock()
b.nextID++
sub := &Subscription{id: b.nextID, handler: h}
b.subs[topic] = append(b.subs[topic], sub)
return sub
}
func (b *EventBus) Unsubscribe(topic string, sub *Subscription) {
b.mu.Lock()
defer b.mu.Unlock()
subs := b.subs[topic]
for i, s := range subs {
if s.id == sub.id {
b.subs[topic] = append(subs[:i], subs[i+1:]...)
return
}
}
}
func (b *EventBus) Publish(topic string, e Event) {
b.mu.RLock()
subs := make([]*Subscription, len(b.subs[topic]))
copy(subs, b.subs[topic])
b.mu.RUnlock()
for _, s := range subs {
s.handler(e)
}
}
func main() {
bus := NewEventBus()
sub := bus.Subscribe("order", func(e Event) {
fmt.Println("观察者 A 收到:", e.Topic)
})
bus.Subscribe("order", func(e Event) {
fmt.Println("观察者 B 收到:", e.Topic)
})
bus.Publish("order", Event{Topic: "order.created"})
bus.Unsubscribe("order", sub)
fmt.Println("--- 取消 A 后 ---")
bus.Publish("order", Event{Topic: "order.paid"})
}三、Go 实现:channel 实现观察者
Go 的 channel 天然适合发布-订阅:每个观察者持有一个 channel,Subject 把事件推到 channel 里,观察者各自从 channel 读。
go
package main
import (
"fmt"
"sync"
)
type Event struct {
Topic string
Payload any
}
type ChannelSubject struct {
mu sync.Mutex
subs []chan Event
closed bool
}
func NewChannelSubject() *ChannelSubject {
return &ChannelSubject{}
}
// Subscribe 返回一个 channel,调用方从中读事件
func (s *ChannelSubject) Subscribe() <-chan Event {
ch := make(chan Event, 16) // 带缓冲,避免慢观察者阻塞发布
s.mu.Lock()
defer s.mu.Unlock()
s.subs = append(s.subs, ch)
return ch
}
func (s *ChannelSubject) Unsubscribe(ch <-chan Event) {
s.mu.Lock()
defer s.mu.Unlock()
for i, c := range s.subs {
if c == ch {
close(s.subs[i])
s.subs = append(s.subs[:i], s.subs[i+1:]...)
return
}
}
}
func (s *ChannelSubject) Publish(e Event) {
s.mu.Lock()
if s.closed {
s.mu.Unlock()
return
}
for _, ch := range s.subs {
select {
case ch <- e:
default:
// 缓冲满则丢弃,避免阻塞(也可改成阻塞)
}
}
s.mu.Unlock()
}
func (s *ChannelSubject) Close() {
s.mu.Lock()
defer s.mu.Unlock()
s.closed = true
for _, ch := range s.subs {
close(ch)
}
s.subs = nil
}
func main() {
sub := NewChannelSubject()
var wg sync.WaitGroup
// 两个观察者,各自从 channel 读
for _, name := range []string{"A", "B"} {
wg.Add(1)
ch := sub.Subscribe()
go func(name string, ch <-chan Event) {
defer wg.Done()
for e := range ch {
fmt.Printf("[%s] 收到: %s\n", name, e.Topic)
}
fmt.Printf("[%s] channel 关闭,退出\n", name)
}(name, ch)
}
sub.Publish(Event{Topic: "e1"})
sub.Publish(Event{Topic: "e2"})
sub.Close()
wg.Wait()
}channel 方式的特点:
- 解耦更强:发布者和订阅者通过 channel 通信,互不知道对方。
- 天然异步:订阅者在自己 goroutine 里读,互不阻塞。
- 背压控制:通过缓冲大小、
select default控制慢订阅者的影响。 - 可组合 select:订阅者可以用
select同时等待多个 channel(事件源、超时、取消)。
四、channel vs 回调对比
| 维度 | 回调 | channel |
|---|---|---|
| 解耦度 | 中(持有函数引用) | 高(只通过 channel) |
| 并发模型 | 同步或异步需自己控制 | 天然异步 |
| 错误传播 | 回调返回 error 易处理 | channel 通常不返回 error,需封装 |
| 背压 | 难(同步会阻塞,异步需自己限流) | 易(缓冲 + select) |
| 取消机制 | 需自己实现 | close(channel) 或 context |
| 生命周期 | 需手动 Unsubscribe | close 即可 |
| 复杂度 | 简单 | 略复杂 |
| 适合场景 | 少量观察者、同步通知 | 多观察者、异步、流式 |
经验:
- 简单的「事件钩子」用回调。
- 复杂的「事件流」用 channel。
- 需要与
select、context、超时组合时,必用 channel。
五、观察者模式与 Go 的 select 多路复用
channel 实现的观察者最大的优势是可以和 select 配合,让订阅者同时等待多个事件源:
go
package main
import (
"context"
"fmt"
"time"
)
func main() {
orders := make(chan string)
payments := make(chan string)
ctx, cancel := context.WithTimeout(context.Background(), 300*time.Millisecond)
defer cancel()
// 模拟事件源
go func() {
time.Sleep(50 * time.Millisecond)
orders <- "ORD-001"
time.Sleep(50 * time.Millisecond)
payments <- "PAY-001"
}()
// 订阅者用 select 同时等订单、支付、超时
for {
select {
case o := <-orders:
fmt.Println("新订单:", o)
case p := <-payments:
fmt.Println("新支付:", p)
case <-ctx.Done():
fmt.Println("超时退出:", ctx.Err())
return
}
}
}这种能力是回调实现做不到的——回调是「推」模型,订阅者被动接受;channel + select 是「拉」模型,订阅者主动选择处理哪个事件。Go 的并发哲学高度依赖这种模式。
六、实战示例:事件总线(Event Bus)
下面写一个稍微完整的事件总线,支持按 topic 订阅、异步分发、优雅关闭:
go
package main
import (
"context"
"fmt"
"sync"
"time"
)
type Event struct {
Topic string
Payload any
}
type Handler func(Event)
type EventBus struct {
mu sync.RWMutex
subscribers map[string][]chan Event
bufferSize int
closed bool
}
func NewEventBus(bufferSize int) *EventBus {
return &EventBus{
subscribers: map[string][]chan Event{},
bufferSize: bufferSize,
}
}
func (b *EventBus) Subscribe(topic string) <-chan Event {
b.mu.Lock()
defer b.mu.Unlock()
ch := make(chan Event, b.bufferSize)
b.subscribers[topic] = append(b.subscribers[topic], ch)
return ch
}
// SubscribeWithHandler 启动一个 goroutine 自动消费
func (b *EventBus) SubscribeWithHandler(ctx context.Context, topic string, h Handler) {
ch := b.Subscribe(topic)
go func() {
for {
select {
case <-ctx.Done():
return
case e, ok := <-ch:
if !ok {
return
}
h(e)
}
}
}()
}
func (b *EventBus) Publish(e Event) {
b.mu.RLock()
defer b.mu.RUnlock()
if b.closed {
return
}
for _, ch := range b.subscribers[e.Topic] {
select {
case ch <- e:
default:
fmt.Printf("[DROP] topic=%s 缓冲满,丢弃\n", e.Topic)
}
}
}
func (b *EventBus) Close() {
b.mu.Lock()
defer b.mu.Unlock()
b.closed = true
for _, subs := range b.subscribers {
for _, ch := range subs {
close(ch)
}
}
b.subscribers = map[string][]chan Event{}
}
func main() {
bus := NewEventBus(8)
ctx, cancel := context.WithTimeout(context.Background(), 500*time.Millisecond)
defer cancel()
// 订阅 order.created
bus.SubscribeWithHandler(ctx, "order.created", func(e Event) {
fmt.Printf("[日志] %s: %v\n", e.Topic, e.Payload)
})
bus.SubscribeWithHandler(ctx, "order.created", func(e Event) {
fmt.Printf("[统计] %s: %v\n", e.Topic, e.Payload)
})
bus.SubscribeWithHandler(ctx, "order.paid", func(e Event) {
fmt.Printf("[发货] %s: %v\n", e.Topic, e.Payload)
})
// 发布事件
bus.Publish(Event{Topic: "order.created", Payload: "ORD-001"})
bus.Publish(Event{Topic: "order.paid", Payload: "ORD-001"})
bus.Publish(Event{Topic: "order.created", Payload: "ORD-002"})
// 等待处理完
time.Sleep(100 * time.Millisecond)
bus.Close()
cancel()
time.Sleep(50 * time.Millisecond)
fmt.Println("结束")
}这个事件总线综合体现了:
- 按 topic 分组订阅。
- 用
context支持取消。 - 缓冲 +
select default防止慢订阅者阻塞发布。 Close优雅关闭所有 channel。
七、标准库中的观察者:context 的取消传播
Go 标准库里没有叫「Observer」的类型,但 context 包的取消传播就是观察者模式的应用。
context.Context 提供 Done() <-chan struct{},所有监听这个 context 的 goroutine 都在「订阅」取消事件。当调用 cancel() 时,所有 Done() channel 被关闭,所有订阅者立刻收到通知。
go
package main
import (
"context"
"fmt"
"time"
)
func worker(ctx context.Context, name string) {
for {
select {
case <-ctx.Done():
fmt.Printf("[%s] 收到取消,退出: %v\n", name, ctx.Err())
return
default:
fmt.Printf("[%s] 工作中...\n", name)
time.Sleep(50 * time.Millisecond)
}
}
}
func main() {
ctx, cancel := context.WithCancel(context.Background())
// 多个 worker 都「订阅」了 ctx 的取消事件
go worker(ctx, "A")
go worker(ctx, "B")
go worker(ctx, "C")
time.Sleep(150 * time.Millisecond)
cancel() // 发布「取消」事件
time.Sleep(50 * time.Millisecond)
}context 的设计完美体现了观察者模式:cancel() 是发布者,所有 ctx.Done() 的监听者是订阅者,发布-订阅通过 channel 完成。这也是 Go 把「并发控制」和「设计模式」融合得最自然的一个例子。
八、使用观察者模式的注意事项
- 避免内存泄漏:订阅者不再需要时必须 Unsubscribe 或关闭 channel,否则 Subject 会一直持有引用。Go 中常用
context配合select <-ctx.Done()来管理订阅者生命周期。 - 避免死锁:同步通知时如果在持有锁的情况下调用回调,回调又去订阅/取消订阅,就会死锁。解决:通知前先复制一份订阅者列表,释放锁后再通知。
- 错误处理:回调实现的观察者可以让 handler 返回 error;channel 实现通常不返回 error,需要单独的 error channel 或封装 Event 结构携带错误。
- 顺序保证:观察者通常不保证顺序。如果业务依赖顺序,要么用单一观察者串行处理,要么显式编号。
- 背压:慢观察者会拖累快发布者。channel + 缓冲 +
select default是 Go 的标准解法。 - 重入问题:观察者的回调里又触发新的事件,可能导致无限递归。需要用状态标记防止。
九、小结
- 观察者模式建立对象间一对多依赖,状态变化时自动通知所有依赖者,又称发布-订阅。
- Go 有两种主流实现:回调(Subject 持有
[]func(Event))和 channel(每个订阅者持有一个 channel)。 - 回调实现简单、适合同步通知;channel 实现解耦更强、天然异步、可与
select/context组合。 - channel +
select是 Go 独有的优势:订阅者可以同时等待多个事件源、超时、取消,这是回调实现做不到的。 context包的取消传播是标准库中观察者模式的经典应用:cancel()是发布者,所有ctx.Done()监听者是订阅者。- 实战中要注意内存泄漏(及时取消订阅)、死锁(通知前复制列表)、背压(缓冲 + select default)。
- 完整的事件总线通常包含:按 topic 订阅、异步分发、context 取消、优雅关闭。
下一篇讲策略模式,看 Go 如何用「函数策略」让代码比 Java 简洁一个量级。