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.

Software Architecture & EngineeringCreated Updated 1 min readhistorical

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

6. References

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub