NOTE
4.4 发布订阅模式
1. 发布订阅模式是什么 - 有三个角色:发布者、事件中心、订阅者 - 订阅者需要向事件中心订阅指定的事件 - 发布者向事件中心发布指定事件内容 - 事件中心通知订阅者 - 订阅者收到消息 2. 为什么需要发布订阅模式 - 解耦 - 发布者只关注生产数据,不关心订阅者怎么消费 - 订阅者只关注消费数
这是历史学习笔记,可能存在过时或不完整的理解。
1. 发布订阅模式是什么
- 有三个角色:发布者、事件中心、订阅者
- 订阅者需要向事件中心订阅指定的事件
- 发布者向事件中心发布指定事件内容
- 事件中心通知订阅者
- 订阅者收到消息
2. 为什么需要发布订阅模式
- 解耦
- 发布者只关注生产数据,不关心订阅者怎么消费
- 订阅者只关注消费数据,不关心发布者怎么生产
- 异步
3. 生产者消费者 vs 发布订阅 vs 观察者
- 生产者消费者:生产者生产的数据只能有一个消费者完整消费;如果由多个消费者,那么每个消费者消费一部分数据
- 后两者:生产者生产的数据可以有多个消费者完整消费
- 发布订阅:发布者不需要手动通知订阅者,支持消息路由
- 观察者:目标需要手动通知观察者,不支持消息路由
4. 实现
type Data struct {
topic string
data interface{}
}
func (d Data) String() string {
return fmt.Sprintf("[topic=%v, data=%v]", d.topic, d.data)
}
type IEventCenter interface {
Publish(topic string, data interface{})
Subscribe(topic string, subscriber ISubscriber)
}
type EventCenter struct {
Queue chan *Data
Subscribers map[string][]ISubscriber
}
func (e *EventCenter) Subscribe(topic string, subscriber ISubscriber) {
subscribers, ok := e.Subscribers[topic]
if ok {
subscribers = append(subscribers, subscriber)
} else {
e.Subscribers[topic] = []ISubscriber{subscriber}
}
}
func (e *EventCenter) Publish(topic string, data interface{}) {
subscribers, ok := e.Subscribers[topic]
if !ok {
return
}
group, newCtx := errgroup.WithContext(context.Background())
for _, subscriber := range subscribers {
group.Go(func() error {
subscriber.Notify(newCtx, &Data{
topic: topic,
data: data,
})
return nil
})
}
group.Wait()
}
// IPublisher ...
type IPublisher interface {
Publish(topic string, data interface{})
}
type PublisherImpl struct {
Queue chan *Data
EventCenter IEventCenter
}
func NewPublisherImpl(eventCenter IEventCenter, size int) *PublisherImpl {
p := &PublisherImpl{EventCenter: eventCenter, Queue: make(chan *Data, size)}
p.monitorQueue()
return p
}
func (p *PublisherImpl) Publish(topic string, data interface{}) {
p.Queue <- &Data{
topic: topic,
data: data,
}
}
func (p *PublisherImpl) monitorQueue() {
go func() {
for {
select {
case data := <-p.Queue:
p.EventCenter.Publish(data.topic, data.data)
}
}
}()
}
// ISubscriber ...
type ISubscriber interface {
Notify(ctx context.Context, data *Data)
}
type SubscriberImpl struct {
name string
}
func (s *SubscriberImpl) Notify(ctx context.Context, data *Data) {
fmt.Println(s.name, data)
}
const (
updateTopic = "updateTopic"
deleteTopic = "deleteTopic"
)
func TestPs3(t *testing.T) {
e := &EventCenter{Subscribers: make(map[string][]ISubscriber, 0), Queue: make(chan *Data, 1000)}
updater := &SubscriberImpl{name: "updater"}
deleter := &SubscriberImpl{name: "deleter"}
p := NewPublisherImpl(e, 1000)
e.Subscribe(updateTopic, updater)
e.Subscribe(deleteTopic, deleter)
p.Publish(updateTopic, 11111)
p.Publish(deleteTopic, "删除数据")
p.Publish(deleteTopic, "删除数据2")
time.Sleep(time.Minute)
}
5. 典型应用
消息队列介绍.md(关联笔记尚未公开)