NOTE
Publish–Subscribe Pattern
What the publish–subscribe pattern is, why it is needed, how it differs from producer–consumer and observer patterns, and a Go implementation.
This is a historical learning note and may contain outdated or incomplete understanding.
1. What Is the Publish–Subscribe Pattern?
- There are three roles: publisher, event center, and subscriber.
- A subscriber subscribes to a specified event through the event center.
- A publisher publishes the content of a specified event to the event center.
- The event center notifies subscribers.
- Subscribers receive the message.
2. Why Is the Publish–Subscribe Pattern Needed?
- Decoupling
- The publisher only focuses on producing data and does not care how subscribers consume it.
- The subscriber only focuses on consuming data and does not care how publishers produce it.
- Asynchrony
3. Producer–Consumer vs. Publish–Subscribe vs. Observer
- Producer–consumer: data produced by a producer can be consumed completely by only one consumer; if there are multiple consumers, each consumer consumes part of the data.
- The latter two: data produced by a producer can be consumed completely by multiple consumers.
- Publish–subscribe: the publisher does not need to notify subscribers manually and supports message routing.
- Observer: the subject needs to notify observers manually and does not support message routing.
4. Implementation
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, "delete data")
p.Publish(deleteTopic, "delete data 2")
time.Sleep(time.Minute)
}
5. Typical Application
Introduction to Message Queues.md
Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub