MidТеория4 min

Observer

Паттерн Observer в Go: подписка на события, channel-based observer и event emitter

Observer (Наблюдатель)

Проблема

Один объект (subject) меняет состояние, и несколько других объектов (observers) должны узнать об этом. Polling неэффективен. Observer позволяет подписчикам получать автоматические уведомления при изменениях.

В Go Observer можно реализовать через интерфейсы (классический) или каналы (Go-идиоматичный).

Диаграмма

    Subject (StockPrice)
    +------------------+
    | observers []Obs  |
    | Subscribe(obs)   |
    | Notify()         |
    +------------------+
            |
            +---> Observer A (Dashboard)
            +---> Observer B (AlertSystem)
            +---> Observer C (Logger)

Реализация: Интерфейсный Observer

package stock

import (
    "fmt"
    "sync"
)

// Event represents a price change event.
type Event struct {
    Symbol   string
    Price    float64
    Change   float64 // percentage
    Volume   int64
}

// Observer receives notifications.
type Observer interface {
    OnPriceChange(event Event)
}

// Subject manages observers and sends notifications.
type Subject struct {
    mu        sync.RWMutex
    observers map[string]Observer
}

// NewSubject creates an observable subject.
func NewSubject() *Subject {
    return &Subject{
        observers: make(map[string]Observer),
    }
}

// Subscribe adds an observer.
func (s *Subject) Subscribe(id string, obs Observer) {
    s.mu.Lock()
    defer s.mu.Unlock()
    s.observers[id] = obs
}

// Unsubscribe removes an observer.
func (s *Subject) Unsubscribe(id string) {
    s.mu.Lock()
    defer s.mu.Unlock()
    delete(s.observers, id)
}

// Notify sends an event to all observers.
func (s *Subject) Notify(event Event) {
    s.mu.RLock()
    defer s.mu.RUnlock()
    for _, obs := range s.observers {
        obs.OnPriceChange(event)
    }
}

// --- Stock Ticker (concrete subject) ---

type Ticker struct {
    Subject
    prices map[string]float64
}

func NewTicker() *Ticker {
    return &Ticker{
        Subject: *NewSubject(),
        prices:  make(map[string]float64),
    }
}

func (t *Ticker) UpdatePrice(symbol string, price float64) {
    old := t.prices[symbol]
    t.prices[symbol] = price

    change := 0.0
    if old > 0 {
        change = (price - old) / old * 100
    }

    t.Notify(Event{
        Symbol: symbol,
        Price:  price,
        Change: change,
    })
}

Конкретные Observer

// Dashboard displays real-time prices.
type Dashboard struct{}

func (d *Dashboard) OnPriceChange(e Event) {
    fmt.Printf("[Dashboard] %s: $%.2f (%+.1f%%)\n", e.Symbol, e.Price, e.Change)
}

// AlertSystem sends alerts on big changes.
type AlertSystem struct {
    threshold float64 // percentage
}

func NewAlertSystem(threshold float64) *AlertSystem {
    return &AlertSystem{threshold: threshold}
}

func (a *AlertSystem) OnPriceChange(e Event) {
    if e.Change > a.threshold || e.Change < -a.threshold {
        fmt.Printf("[ALERT] %s moved %+.1f%% to $%.2f!\n", e.Symbol, e.Change, e.Price)
    }
}

Go-идиоматичный Observer: каналы

package pubsub

import (
    "context"
    "sync"
)

// Broker is a channel-based observer (pub/sub).
type Broker[T any] struct {
    mu          sync.RWMutex
    subscribers map[string]chan T
    bufSize     int
}

// NewBroker creates a pub/sub broker.
func NewBroker[T any](bufSize int) *Broker[T] {
    return &Broker[T]{
        subscribers: make(map[string]chan T),
        bufSize:     bufSize,
    }
}

// Subscribe returns a channel that receives published events.
func (b *Broker[T]) Subscribe(id string) <-chan T {
    b.mu.Lock()
    defer b.mu.Unlock()

    ch := make(chan T, b.bufSize)
    b.subscribers[id] = ch
    return ch
}

// Unsubscribe removes a subscriber and closes their channel.
func (b *Broker[T]) Unsubscribe(id string) {
    b.mu.Lock()
    defer b.mu.Unlock()

    if ch, ok := b.subscribers[id]; ok {
        close(ch)
        delete(b.subscribers, id)
    }
}

// Publish sends an event to all subscribers.
// Non-blocking: skips slow subscribers.
func (b *Broker[T]) Publish(event T) {
    b.mu.RLock()
    defer b.mu.RUnlock()

    for _, ch := range b.subscribers {
        select {
        case ch <- event:
        default:
            // subscriber is slow, skip
        }
    }
}

// Close unsubscribes all and closes all channels.
func (b *Broker[T]) Close() {
    b.mu.Lock()
    defer b.mu.Unlock()

    for id, ch := range b.subscribers {
        close(ch)
        delete(b.subscribers, id)
    }
}

Использование channel-based observer

func main() {
    broker := pubsub.NewBroker[stock.Event](100)

    // Subscriber 1: Dashboard
    dashCh := broker.Subscribe("dashboard")
    go func() {
        for event := range dashCh {
            fmt.Printf("[Dashboard] %s: $%.2f\n", event.Symbol, event.Price)
        }
    }()

    // Subscriber 2: Alert system
    alertCh := broker.Subscribe("alerts")
    go func() {
        for event := range alertCh {
            if event.Change > 5 || event.Change < -5 {
                fmt.Printf("[ALERT] %s: %+.1f%%\n", event.Symbol, event.Change)
            }
        }
    }()

    // Publisher
    broker.Publish(stock.Event{Symbol: "GOOG", Price: 150.0, Change: 2.5})
    broker.Publish(stock.Event{Symbol: "GOOG", Price: 160.0, Change: 6.7})

    // Cleanup
    broker.Unsubscribe("dashboard")
    broker.Close()
}

Typed Event Bus с фильтрацией

package event

import "context"

// HandlerFunc handles events of a specific type.
type HandlerFunc[T any] func(ctx context.Context, event T)

// FilterFunc decides if an event should be delivered.
type FilterFunc[T any] func(event T) bool

// Subscription wraps a handler with an optional filter.
type Subscription[T any] struct {
    handler HandlerFunc[T]
    filter  FilterFunc[T]
}

// Emitter is a typed event emitter.
type Emitter[T any] struct {
    mu   sync.RWMutex
    subs []Subscription[T]
}

// On subscribes to all events.
func (e *Emitter[T]) On(handler HandlerFunc[T]) {
    e.mu.Lock()
    defer e.mu.Unlock()
    e.subs = append(e.subs, Subscription[T]{handler: handler})
}

// OnFiltered subscribes with a filter.
func (e *Emitter[T]) OnFiltered(filter FilterFunc[T], handler HandlerFunc[T]) {
    e.mu.Lock()
    defer e.mu.Unlock()
    e.subs = append(e.subs, Subscription[T]{handler: handler, filter: filter})
}

// Emit sends an event to matching subscribers.
func (e *Emitter[T]) Emit(ctx context.Context, event T) {
    e.mu.RLock()
    defer e.mu.RUnlock()

    for _, sub := range e.subs {
        if sub.filter != nil && !sub.filter(event) {
            continue
        }
        sub.handler(ctx, event)
    }
}

Когда использовать

Используйте, когда:

  • Изменение одного объекта должно отразиться на других
  • Количество подписчиков неизвестно или меняется
  • Event-driven архитектура
  • Слабая связанность между компонентами

Не используйте, когда:

  • Один подписчик (просто вызовите callback)
  • Порядок уведомления критичен (Observer не гарантирует порядок)
  • Нужна двусторонняя связь (используйте Mediator)

Сравнение: интерфейсный vs channel-based

Критерий Интерфейсный Channel-based
Идиоматичность Средняя Высокая (Go way)
Concurrency Нужен sync Встроена в каналы
Backpressure Нет Да (буфер канала)
Отписка Явная (Unsubscribe) close(channel)
Лучше для Синхронной обработки Асинхронной обработки

Проверь себя

Какой Go-идиоматичный способ реализации Observer?

Почему в Publish() используется select с default?