NOTE

4.4 发布订阅模式

1. 发布订阅模式是什么 - 有三个角色:发布者、事件中心、订阅者 - 订阅者需要向事件中心订阅指定的事件 - 发布者向事件中心发布指定事件内容 - 事件中心通知订阅者 - 订阅者收到消息 2. 为什么需要发布订阅模式 - 解耦 - 发布者只关注生产数据,不关心订阅者怎么消费 - 订阅者只关注消费数

Software Architecture & Engineering创建于 更新于 historical

这是历史学习笔记,可能存在过时或不完整的理解。

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(关联笔记尚未公开)

6. 参考