Skip to content

观察者模式

观察者模式(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
生命周期需手动 Unsubscribeclose 即可
复杂度简单略复杂
适合场景少量观察者、同步通知多观察者、异步、流式

经验:

  • 简单的「事件钩子」用回调。
  • 复杂的「事件流」用 channel。
  • 需要与 selectcontext、超时组合时,必用 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 简洁一个量级。