Observer Pattern #
A new order comes into the e-commerce system. As a result, many things must happen: inventory must be decremented, a confirmation email must be sent, the fulfillment team must be notified, analytics must be updated, and loyalty points must be added to the buyer’s account. Without the Observer Pattern, OrderService has to call all of them explicitly — every new subsystem that needs to know about new orders means one more line of code in OrderService. The OrderService, which should only handle order logic, ends up handling inventory, email, analytics, and loyalty all at once — and every time a new subsystem appears, it has to change. The Observer Pattern reverses this responsibility: OrderService merely announces “a new order has arrived”, and anyone who cares can register to hear that announcement — without OrderService needing to know who is listening or how many.
What Is the Observer Pattern? #
The Observer Pattern is a behavioral design pattern that defines a subscription mechanism to notify many objects about events that occur on the object they are watching, without the watched object (the Subject) needing to know who is watching it or what they do with the information.
Two main roles in the Observer Pattern:
- Subject (Publisher) — the object that holds state and notifies observers when the state changes; it does not know who its observers are, only that they exist
- Observer (Subscriber) — the object that registers itself to receive notifications; it is responsible for what happens when a notification arrives
Three properties define the Observer Pattern:
- Loose coupling — the Subject does not depend on concrete Observers; the two only know each other through an interface
- Dynamic subscription — Observers can be registered and removed at runtime, without the Subject needing to change
- Broadcast — one notification from the Subject reaches all registered Observers at once
sequenceDiagram
participant OS as OrderService (Subject)
participant IS as InventoryObserver
participant ES as EmailObserver
participant AS as AnalyticsObserver
Note over IS,AS: Register themselves with OrderService
IS->>OS: Register(inventoryObs)
ES->>OS: Register(emailObs)
AS->>OS: Register(analyticsObs)
Note over OS: A new order arrives
OS->>OS: CreateOrder(req)
OS->>IS: Notify(OrderCreatedEvent)
OS->>ES: Notify(OrderCreatedEvent)
OS->>AS: Notify(OrderCreatedEvent)
IS-->>IS: decrement stock
ES-->>ES: send confirmation email
AS-->>AS: record eventThe Problem It Solves #
The Observer Pattern solves two interrelated problems: tight coupling between the Subject and the parties reacting to its changes, and the Open-Closed Principle being violated every time a new reaction needs to be added.
The Problem: The Subject Knows Too Much #
// ANTI-PATTERN: OrderService calls every subsystem directly
func (s *OrderService) CreateOrder(req CreateOrderRequest) (*Order, error) {
order, err := s.repo.Create(req)
if err != nil {
return nil, err
}
// OrderService has to know about inventory
s.inventorySvc.DeductStock(order.Items)
// OrderService has to know about email
s.emailSvc.SendConfirmation(order.UserEmail, order)
// OrderService has to know about analytics
s.analyticsSvc.TrackPurchase(order)
// OrderService has to know about loyalty
s.loyaltySvc.AddPoints(order.UserID, order.TotalAmount)
// Add WhatsApp notifications? Modify OrderService again.
// Add fraud detection? Modify OrderService again.
// Total coupling: OrderService knows 4+ subsystems
return order, nil
}
// CORRECT: OrderService only knows that an event exists, not who needs to know
func (s *OrderService) CreateOrder(req CreateOrderRequest) (*Order, error) {
order, err := s.repo.Create(req)
if err != nil {
return nil, err
}
// One line — does not care who is listening or how many
s.eventBus.Publish(OrderCreatedEvent{Order: order})
return order, nil
}
// Add WhatsApp? Register a WhatsAppObserver — OrderService does not change
// Add fraud detection? Register a FraudObserver — OrderService does not change
Structure and Components #
classDiagram
class Subject {
<<interface>>
+Register(observer Observer)
+Unregister(observer Observer)
+Notify(event Event)
}
class Observer {
<<interface>>
+OnEvent(event Event)
}
class OrderEventBus {
-observers map[EventType][]Observer
-mu sync.RWMutex
+Register(eventType EventType, observer Observer)
+Unregister(eventType EventType, observer Observer)
+Publish(event Event)
}
class InventoryObserver {
-inventorySvc InventoryService
+OnEvent(event Event)
}
class EmailObserver {
-emailSvc EmailService
+OnEvent(event Event)
}
class AnalyticsObserver {
-analyticsSvc AnalyticsService
+OnEvent(event Event)
}
Subject <|.. OrderEventBus
Observer <|.. InventoryObserver
Observer <|.. EmailObserver
Observer <|.. AnalyticsObserver
OrderEventBus o-- Observer : manages list| Component | Role | Characteristics |
|---|---|---|
| Subject interface | Contract for registering and notifying observers | Register, Unregister, Notify |
| Concrete Subject | Stores the observer list, triggers notifications | Thread-safe with a mutex |
| Observer interface | Contract every subscriber must satisfy | Usually one method: OnEvent or Update |
| Concrete Observer | Reacts to events according to its responsibility | Independent of each other |
Full Implementation: Order Event System #
Event Types and Interfaces #
package events
import "time"
// EventType defines the kinds of events that can be published.
type EventType string
const (
EventOrderCreated EventType = "order.created"
EventOrderPaid EventType = "order.paid"
EventOrderShipped EventType = "order.shipped"
EventOrderCancelled EventType = "order.cancelled"
)
// Event is the interface for every event in the system.
type Event interface {
Type() EventType
OccurredAt() time.Time
}
// OrderCreatedEvent is published when a new order is successfully created.
type OrderCreatedEvent struct {
OrderID string
UserID string
UserEmail string
Items []OrderItem
TotalAmount float64
CreatedAt time.Time
}
func (e OrderCreatedEvent) Type() EventType { return EventOrderCreated }
func (e OrderCreatedEvent) OccurredAt() time.Time { return e.CreatedAt }
// OrderPaidEvent is published when a payment is successfully confirmed.
type OrderPaidEvent struct {
OrderID string
UserID string
UserEmail string
TransactionID string
Amount float64
PaidAt time.Time
}
func (e OrderPaidEvent) Type() EventType { return EventOrderPaid }
func (e OrderPaidEvent) OccurredAt() time.Time { return e.PaidAt }
// OrderShippedEvent is published when an order is shipped.
type OrderShippedEvent struct {
OrderID string
UserID string
UserEmail string
TrackingNumber string
Courier string
ShippedAt time.Time
}
func (e OrderShippedEvent) Type() EventType { return EventOrderShipped }
func (e OrderShippedEvent) OccurredAt() time.Time { return e.ShippedAt }
// OrderItem is an item within an order.
type OrderItem struct {
ProductID string
Quantity int
Price float64
}
// Observer is the interface every subscriber must implement.
type Observer interface {
OnEvent(event Event)
ObserverName() string // for logging and debugging
}
Concrete Subject: Thread-Safe EventBus #
package events
import (
"fmt"
"log/slog"
"sync"
)
// EventBus is the Subject — it manages the observer list per event type
// and distributes events to the relevant observers.
type EventBus struct {
mu sync.RWMutex
observers map[EventType][]Observer
logger *slog.Logger
}
// NewEventBus creates a new, ready-to-use EventBus.
func NewEventBus(logger *slog.Logger) *EventBus {
return &EventBus{
observers: make(map[EventType][]Observer),
logger: logger,
}
}
// Register subscribes an observer to a specific event type.
// The same observer can be registered for several event types.
func (b *EventBus) Register(eventType EventType, observer Observer) {
b.mu.Lock()
defer b.mu.Unlock()
// Check whether the observer is already registered for this event type
for _, existing := range b.observers[eventType] {
if existing == observer {
b.logger.Warn("observer already registered",
"event_type", eventType,
"observer", observer.ObserverName())
return
}
}
b.observers[eventType] = append(b.observers[eventType], observer)
b.logger.Info("observer registered",
"event_type", eventType,
"observer", observer.ObserverName(),
"total", len(b.observers[eventType]),
)
}
// Unregister removes an observer from a specific event type.
func (b *EventBus) Unregister(eventType EventType, observer Observer) {
b.mu.Lock()
defer b.mu.Unlock()
observers := b.observers[eventType]
for i, existing := range observers {
if existing == observer {
b.observers[eventType] = append(observers[:i], observers[i+1:]...)
b.logger.Info("observer unregistered",
"event_type", eventType,
"observer", observer.ObserverName(),
)
return
}
}
}
// Publish distributes an event to all registered observers.
// Notification is synchronous — all observers finish before Publish returns.
func (b *EventBus) Publish(event Event) {
b.mu.RLock()
observers := make([]Observer, len(b.observers[event.Type()]))
copy(observers, b.observers[event.Type()])
b.mu.RUnlock()
b.logger.Info("publishing event",
"event_type", event.Type(),
"observer_count", len(observers),
"occurred_at", event.OccurredAt(),
)
for _, observer := range observers {
func() {
defer func() {
if r := recover(); r != nil {
b.logger.Error("observer panicked",
"observer", observer.ObserverName(),
"event_type", event.Type(),
"panic", r,
)
}
}()
observer.OnEvent(event)
}()
}
}
// PublishAsync distributes an event asynchronously via goroutines.
// Use it for observers that can run in parallel and independently.
func (b *EventBus) PublishAsync(event Event) {
b.mu.RLock()
observers := make([]Observer, len(b.observers[event.Type()]))
copy(observers, b.observers[event.Type()])
b.mu.RUnlock()
for _, observer := range observers {
go func(obs Observer) {
defer func() {
if r := recover(); r != nil {
b.logger.Error("async observer panicked",
"observer", obs.ObserverName(),
"event_type", event.Type(),
"panic", fmt.Sprintf("%v", r),
)
}
}()
obs.OnEvent(event)
}(observer)
}
}
// ObserverCount returns the number of observers registered for a specific event type.
func (b *EventBus) ObserverCount(eventType EventType) int {
b.mu.RLock()
defer b.mu.RUnlock()
return len(b.observers[eventType])
}
Concrete Observers #
package observers
import (
"context"
"fmt"
"log/slog"
"myapp/events"
)
// InventoryObserver decrements stock when an order is created.
type InventoryObserver struct {
inventorySvc InventoryService
logger *slog.Logger
}
func NewInventoryObserver(svc InventoryService, logger *slog.Logger) *InventoryObserver {
return &InventoryObserver{inventorySvc: svc, logger: logger}
}
func (o *InventoryObserver) OnEvent(event events.Event) {
e, ok := event.(events.OrderCreatedEvent)
if !ok {
return // irrelevant event, ignore it
}
ctx := context.Background()
if err := o.inventorySvc.DeductStock(ctx, e.Items); err != nil {
o.logger.Error("failed to deduct stock",
"order_id", e.OrderID,
"error", err,
)
return
}
o.logger.Info("stock deducted", "order_id", e.OrderID, "items", len(e.Items))
}
func (o *InventoryObserver) ObserverName() string { return "InventoryObserver" }
// EmailObserver sends confirmation emails when orders are created and paid.
type EmailObserver struct {
emailSvc EmailService
logger *slog.Logger
}
func NewEmailObserver(svc EmailService, logger *slog.Logger) *EmailObserver {
return &EmailObserver{emailSvc: svc, logger: logger}
}
func (o *EmailObserver) OnEvent(event events.Event) {
ctx := context.Background()
switch e := event.(type) {
case events.OrderCreatedEvent:
if err := o.emailSvc.SendOrderConfirmation(ctx, e.UserEmail, e.OrderID, e.TotalAmount); err != nil {
o.logger.Error("failed to send confirmation email",
"order_id", e.OrderID,
"email", e.UserEmail,
"error", err,
)
}
case events.OrderShippedEvent:
if err := o.emailSvc.SendShippingNotification(ctx, e.UserEmail, e.OrderID, e.TrackingNumber, e.Courier); err != nil {
o.logger.Error("failed to send shipping email",
"order_id", e.OrderID,
"error", err,
)
}
default:
// Other events are not handled by EmailObserver
}
}
func (o *EmailObserver) ObserverName() string { return "EmailObserver" }
// AnalyticsObserver records every event to the analytics platform.
type AnalyticsObserver struct {
analyticsSvc AnalyticsService
logger *slog.Logger
}
func NewAnalyticsObserver(svc AnalyticsService, logger *slog.Logger) *AnalyticsObserver {
return &AnalyticsObserver{analyticsSvc: svc, logger: logger}
}
func (o *AnalyticsObserver) OnEvent(event events.Event) {
ctx := context.Background()
properties := map[string]interface{}{
"event_type": string(event.Type()),
"occurred_at": event.OccurredAt(),
}
switch e := event.(type) {
case events.OrderCreatedEvent:
properties["order_id"] = e.OrderID
properties["user_id"] = e.UserID
properties["amount"] = e.TotalAmount
properties["item_count"] = len(e.Items)
case events.OrderPaidEvent:
properties["order_id"] = e.OrderID
properties["user_id"] = e.UserID
properties["amount"] = e.Amount
properties["transaction_id"] = e.TransactionID
}
if err := o.analyticsSvc.Track(ctx, string(event.Type()), properties); err != nil {
o.logger.Warn("failed to track analytics event", "error", err)
// Non-fatal: analytics failure must not stop the main flow
}
}
func (o *AnalyticsObserver) ObserverName() string { return "AnalyticsObserver" }
// LoyaltyObserver adds loyalty points when a payment succeeds.
type LoyaltyObserver struct {
loyaltySvc LoyaltyService
logger *slog.Logger
}
func NewLoyaltyObserver(svc LoyaltyService, logger *slog.Logger) *LoyaltyObserver {
return &LoyaltyObserver{loyaltySvc: svc, logger: logger}
}
func (o *LoyaltyObserver) OnEvent(event events.Event) {
e, ok := event.(events.OrderPaidEvent)
if !ok {
return
}
// Calculate points: 1 point per 10,000 paid
points := int(e.Amount / 10000)
if points == 0 {
return
}
ctx := context.Background()
if err := o.loyaltySvc.AddPoints(ctx, e.UserID, points); err != nil {
o.logger.Error("failed to add loyalty points",
"user_id", e.UserID,
"points", points,
"error", err,
)
return
}
o.logger.Info("loyalty points added", "user_id", e.UserID, "points", points)
}
func (o *LoyaltyObserver) ObserverName() string { return "LoyaltyObserver" }
Subject: A Clean OrderService #
package order
import (
"context"
"fmt"
"time"
"myapp/events"
)
// OrderService is the Subject — it publishes events, caring nothing about its observers.
type OrderService struct {
repo OrderRepository
eventBus *events.EventBus
}
func NewOrderService(repo OrderRepository, eventBus *events.EventBus) *OrderService {
return &OrderService{repo: repo, eventBus: eventBus}
}
func (s *OrderService) CreateOrder(ctx context.Context, req CreateOrderRequest) (*Order, error) {
order, err := s.repo.Create(ctx, req)
if err != nil {
return nil, fmt.Errorf("failed to create order: %w", err)
}
// Publish the event — does not know or care who subscribes
s.eventBus.Publish(events.OrderCreatedEvent{
OrderID: order.ID,
UserID: req.UserID,
UserEmail: req.UserEmail,
Items: req.Items,
TotalAmount: order.TotalAmount,
CreatedAt: time.Now(),
})
return order, nil
}
func (s *OrderService) ConfirmPayment(ctx context.Context, orderID, transactionID string, amount float64) error {
order, err := s.repo.UpdateStatus(ctx, orderID, "paid")
if err != nil {
return fmt.Errorf("failed to update order status: %w", err)
}
s.eventBus.Publish(events.OrderPaidEvent{
OrderID: orderID,
UserID: order.UserID,
UserEmail: order.UserEmail,
TransactionID: transactionID,
Amount: amount,
PaidAt: time.Now(),
})
return nil
}
func (s *OrderService) ShipOrder(ctx context.Context, orderID, trackingNumber, courier string) error {
order, err := s.repo.UpdateStatus(ctx, orderID, "shipped")
if err != nil {
return fmt.Errorf("failed to update order status: %w", err)
}
s.eventBus.Publish(events.OrderShippedEvent{
OrderID: orderID,
UserID: order.UserID,
UserEmail: order.UserEmail,
TrackingNumber: trackingNumber,
Courier: courier,
ShippedAt: time.Now(),
})
return nil
}
Wiring in main.go #
func main() {
logger := slog.Default()
eventBus := events.NewEventBus(logger)
// Create the observers
inventoryObs := observers.NewInventoryObserver(inventorySvc, logger)
emailObs := observers.NewEmailObserver(emailSvc, logger)
analyticsObs := observers.NewAnalyticsObserver(analyticsSvc, logger)
loyaltyObs := observers.NewLoyaltyObserver(loyaltySvc, logger)
// Register observers for the relevant events
eventBus.Register(events.EventOrderCreated, inventoryObs)
eventBus.Register(events.EventOrderCreated, emailObs)
eventBus.Register(events.EventOrderCreated, analyticsObs)
eventBus.Register(events.EventOrderPaid, loyaltyObs)
eventBus.Register(events.EventOrderPaid, analyticsObs)
eventBus.Register(events.EventOrderShipped, emailObs)
eventBus.Register(events.EventOrderShipped, analyticsObs)
// OrderService does not need to know any of this
orderSvc := order.NewOrderService(orderRepo, eventBus)
// Adding a FraudDetectionObserver later:
// eventBus.Register(events.EventOrderCreated, fraudObs)
// — no other code changes
}
Async Observer: Non-Blocking Notification #
For observers that can run independently and do not need to wait for each other, async notification via goroutines is more appropriate.
// AsyncEventBus distributes an event to each observer in a separate goroutine.
// Use it when observers do not need a guaranteed execution order.
type AsyncEventBus struct {
EventBus
errorHandler func(obs Observer, event Event, err interface{})
}
func NewAsyncEventBus(logger *slog.Logger) *AsyncEventBus {
bus := &AsyncEventBus{}
bus.observers = make(map[EventType][]Observer)
bus.logger = logger
bus.errorHandler = func(obs Observer, event Event, err interface{}) {
logger.Error("async observer error",
"observer", obs.ObserverName(),
"event_type", event.Type(),
"error", fmt.Sprintf("%v", err),
)
}
return bus
}
func (b *AsyncEventBus) Publish(event Event) {
b.mu.RLock()
observers := make([]Observer, len(b.observers[event.Type()]))
copy(observers, b.observers[event.Type()])
b.mu.RUnlock()
var wg sync.WaitGroup
for _, observer := range observers {
wg.Add(1)
go func(obs Observer) {
defer wg.Done()
defer func() {
if r := recover(); r != nil {
b.errorHandler(obs, event, r)
}
}()
obs.OnEvent(event)
}(observer)
}
// Wait for all to finish — or not, depending on the need
// wg.Wait() // uncomment if you need to wait for all observers to finish
}
Sync vs async comparison:
| Aspect | Synchronous Publish | Asynchronous Publish |
|---|---|---|
| Execution order | Guaranteed — observers are called in sequence | Not guaranteed — runs in parallel |
| Error handling | An observer error can delay the next observer | One observer’s error does not affect others |
| Rollback | Rollback can be implemented if an observer fails | Hard — the event is already published |
| Performance | A slow observer blocks the Subject | Observers run in parallel |
| Best for | Interdependent observers that need ordering | Independent, best-effort observers |
Preventing Memory Leaks #
One of the biggest mishandled problems of the Observer Pattern is memory leaks — observers that are no longer active remain registered and prevent garbage collection.
// ANTI-PATTERN: observers are never unregistered
type SessionHandler struct {
eventBus *events.EventBus
}
func (h *SessionHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
// A new observer is registered on every request — never removed!
obs := &RequestScopedObserver{requestID: r.Header.Get("X-Request-ID")}
h.eventBus.Register(events.EventOrderCreated, obs)
// After the request finishes, obs is not removed → it keeps accumulating
}
// CORRECT: always unregister observers that are no longer needed
func (h *SessionHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
obs := &RequestScopedObserver{requestID: r.Header.Get("X-Request-ID")}
h.eventBus.Register(events.EventOrderCreated, obs)
defer h.eventBus.Unregister(events.EventOrderCreated, obs) // automatic cleanup
// handle request...
}
// Or: use a WeakReference-style approach with a timeout
type TimedObserver struct {
observer Observer
expiresAt time.Time
}
// EventBus with auto-cleanup of expired observers
func (b *EventBus) CleanupExpired() {
b.mu.Lock()
defer b.mu.Unlock()
now := time.Now()
for eventType, observers := range b.observers {
var active []Observer
for _, obs := range observers {
if timed, ok := obs.(*TimedObserver); ok {
if now.Before(timed.expiresAt) {
active = append(active, obs)
}
} else {
active = append(active, obs)
}
}
b.observers[eventType] = active
}
}
Observers Living Longer Than Expected
If the Subject holds a reference to an Observer, the Observer will not be garbage collected even when all other references to it are gone. This is a very common source of memory leaks. Always provide an
Unregistermechanism and make sure the Observer lifecycle is managed properly — especially for Observers created per-request or per-session.
Testing the Observer Pattern #
// MockObserver for testing — records every event it receives
type MockObserver struct {
ReceivedEvents []events.Event
mu sync.Mutex
ShouldPanic bool
}
func (m *MockObserver) OnEvent(event events.Event) {
if m.ShouldPanic {
panic("mock panic")
}
m.mu.Lock()
defer m.mu.Unlock()
m.ReceivedEvents = append(m.ReceivedEvents, event)
}
func (m *MockObserver) ObserverName() string { return "MockObserver" }
func (m *MockObserver) EventCount() int {
m.mu.Lock()
defer m.mu.Unlock()
return len(m.ReceivedEvents)
}
func TestEventBus_Publish_NotifiesAllObservers(t *testing.T) {
bus := events.NewEventBus(slog.Default())
obs1 := &MockObserver{}
obs2 := &MockObserver{}
bus.Register(events.EventOrderCreated, obs1)
bus.Register(events.EventOrderCreated, obs2)
event := events.OrderCreatedEvent{
OrderID: "ORD-001", UserID: "user-1",
TotalAmount: 150000, CreatedAt: time.Now(),
}
bus.Publish(event)
if obs1.EventCount() != 1 {
t.Errorf("expected obs1 to receive 1 event, got %d", obs1.EventCount())
}
if obs2.EventCount() != 1 {
t.Errorf("expected obs2 to receive 1 event, got %d", obs2.EventCount())
}
}
func TestEventBus_Unregister_StopsNotification(t *testing.T) {
bus := events.NewEventBus(slog.Default())
obs := &MockObserver{}
bus.Register(events.EventOrderCreated, obs)
bus.Publish(events.OrderCreatedEvent{OrderID: "ORD-001", CreatedAt: time.Now()})
bus.Unregister(events.EventOrderCreated, obs)
bus.Publish(events.OrderCreatedEvent{OrderID: "ORD-002", CreatedAt: time.Now()})
if obs.EventCount() != 1 {
t.Errorf("expected 1 event after unregister, got %d", obs.EventCount())
}
}
func TestEventBus_PanicInObserver_DoesNotStopOthers(t *testing.T) {
bus := events.NewEventBus(slog.Default())
panicObs := &MockObserver{ShouldPanic: true}
normalObs := &MockObserver{}
bus.Register(events.EventOrderCreated, panicObs)
bus.Register(events.EventOrderCreated, normalObs)
// Publish must not panic even when one observer panics
require.NotPanics(t, func() {
bus.Publish(events.OrderCreatedEvent{OrderID: "ORD-001", CreatedAt: time.Now()})
})
// The normal observer still receives the event
if normalObs.EventCount() != 1 {
t.Errorf("normal observer should still receive event despite panic in another observer")
}
}
func TestEventBus_TypeFiltering(t *testing.T) {
bus := events.NewEventBus(slog.Default())
obs := &MockObserver{}
// Only subscribe to EventOrderPaid
bus.Register(events.EventOrderPaid, obs)
// Publish EventOrderCreated — obs should not receive it
bus.Publish(events.OrderCreatedEvent{OrderID: "ORD-001", CreatedAt: time.Now()})
if obs.EventCount() != 0 {
t.Error("observer should not receive event for unsubscribed type")
}
// Publish EventOrderPaid — should be received
bus.Publish(events.OrderPaidEvent{OrderID: "ORD-001", Amount: 100000, PaidAt: time.Now()})
if obs.EventCount() != 1 {
t.Errorf("expected 1 event for subscribed type, got %d", obs.EventCount())
}
}
func TestEventBus_ConcurrentPublish(t *testing.T) {
bus := events.NewEventBus(slog.Default())
obs := &MockObserver{}
bus.Register(events.EventOrderCreated, obs)
const goroutines = 100
var wg sync.WaitGroup
for i := 0; i < goroutines; i++ {
wg.Add(1)
go func(i int) {
defer wg.Done()
bus.Publish(events.OrderCreatedEvent{
OrderID: fmt.Sprintf("ORD-%d", i),
CreatedAt: time.Now(),
})
}(i)
}
wg.Wait()
if obs.EventCount() != goroutines {
t.Errorf("expected %d events from concurrent publish, got %d", goroutines, obs.EventCount())
}
}
Observer vs Event Bus vs Mediator #
The Observer Pattern can be implemented in several ways. It is important to understand the differences.
| Aspect | Classic Observer | Event Bus (Typed) | Mediator |
|---|---|---|---|
| Does the Subject know its Observers? | Yes — it stores a list | No — only knows EventType | No — everything goes through the mediator |
| Do Observers know the Subject? | No — only know the event | No — only know EventType | No |
| Coupling | Subject-Observer loosely coupled | Most loosely coupled | Components do not know each other |
| Event filtering | Observers filter themselves | Filter by type at the bus | The Mediator decides who needs to know |
| Best for | One subject, several observers | Many subjects, many observers | Complex communication between components |
When to Use and When Not to #
USE Observer if:
✓ A change in one object needs to be known by other objects without tight coupling
✓ The number and identity of observers are not known at compile time
✓ Observers can be registered and removed dynamically
✓ You are building an event-driven or reactive system
✓ You want to apply the Open-Closed Principle to reactions to events
AVOID Observer if:
✗ Notification chains get too long and make debugging difficult
✗ Observers need to know about other observers — use a Mediator
✗ Observer order is critical and cannot be guaranteed
✗ The system is very small with 1-2 reactions that will never change
Observer Review Checklist #
DESIGN:
□ The Subject does not depend on concrete Observers — only on the interface
□ Observers can be registered and removed without changing the Subject
□ Events carry enough information that Observers do not need to query the Subject back
□ Each Observer handles only relevant events — ignores the rest
IMPLEMENTATION:
□ The EventBus is thread-safe — Register, Unregister, Publish use a mutex
□ Publish copies the observer list before calling (avoid deadlock)
□ A panic in one observer does not stop notifications to other observers
□ Observer errors are not propagated to the Subject — the Subject does not care
LIFECYCLE:
□ There is an easy-to-use Unregister mechanism
□ Observers are unregistered when no longer needed
□ No circular references between Subject and Observer
TESTING:
□ A MockObserver records all events for assertions
□ Test that all observers are notified
□ Test that Unregister stops notifications
□ Test that a panic in one observer does not affect others
□ Test thread-safety with concurrent Publish
Summary #
- Observer defines a subscription mechanism — the Subject announces changes, Observers react; the two do not need to know each other concretely.
- True loose coupling: the Subject does not know who its observers are, how many, or what they do — it only calls
OnEvent()through an interface.- A typed Event Bus is an evolution of the classic Observer — events are filtered by type at the bus, not in each observer; more scalable for systems with many events.
- Thread safety is mandatory:
Register,Unregister, andPublishare often called from different goroutines; usesync.RWMutexand copy the observer list before iterating.- Panic recovery inside Publish — one panicking observer must not stop notifications to other observers; wrap every observer call with
recover().- Memory leaks are a real risk — unregistered observers prevent garbage collection; always use
defer Unregister()for observers with a limited lifecycle.- Sync vs Async: synchronous for observers that need ordering; async (goroutines) for independent observers that can run in parallel.
- Distinguish it from Mediator: Observer is for “one subject, many subscribers”; Mediator is for “many components communicating through an intermediary”.