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) |
| Лучше для | Синхронной обработки | Асинхронной обработки |