观察者设计模式使订阅者能够向提供程序注册和接收通知。 它适用于需要基于推送的通知的任何方案。 该模式定义提供程序(也称为主体或可观察对象)和零、一个或多个观察者。 观察程序向提供程序注册,每当发生预定义的条件、事件或状态更改时,提供程序都会通过调用委托自动通知所有观察程序。 在此方法调用中,提供程序还可以向观察者提供当前状态信息。 在 .NET 中,通过实现泛型 System.IObservable<T> 和 System.IObserver<T> 接口来应用观察程序设计模式。 泛型类型参数表示提供通知信息的类型。
何时应用模式
观察者设计模式适用于基于推送的分布式通知,因为它支持在两个不同的组件或应用程序层(如数据源(业务逻辑)层和用户界面(显示)层之间实现干净分离。 每当提供程序使用回调向客户端提供当前信息时,都可以实现该模式。
实现模式需要提供以下详细信息:
提供者或主题,即向观察者发送通知的对象。 提供者是实现接口的 IObservable<T> 类或结构体。 提供程序必须实现单个方法, IObservable<T>.Subscribe该方法由希望从提供程序接收通知的观察程序调用。
观察者是一个接收提供程序通知的对象。 观察者是实现 IObserver<T> 接口的类或结构。 观察程序必须实现三种方法,所有这些方法均由提供程序调用:
- IObserver<T>.OnNext,它为观察者提供新的或当前信息。
- IObserver<T>.OnError,它通知观察者发生了错误。
- IObserver<T>.OnCompleted,指示提供程序已完成发送通知。
允许提供者跟踪观察者的机制。 通常,提供程序使用容器对象(如 System.Collections.Generic.List<T> 对象)来保存对 IObserver<T> 已订阅通知的实现的引用。 使用存储容器实现此目的,使提供者可以处理零到无限数量的观察者。 未定义观察程序接收通知的顺序;提供程序可以使用任何方法来确定顺序。
一种 IDisposable 实现,允许提供程序在通知完成后删除观察者。 观察程序从IDisposable该方法接收对Subscribe实现的引用,因此,在提供程序完成发送通知之前,观察程序还可以调用IDisposable.Dispose该方法取消订阅。
一个对象,包含提供者发送到其观察者的数据。 此对象的类型对应于IObservable<T>和IObserver<T>接口的泛型类型参数。 尽管此对象可以与 IObservable<T> 实现相同,但最常见的是单独的类型。
注释
除了实现观察程序设计模式之外,你可能还有兴趣浏览使用 IObservable<T> 和 IObserver<T> 接口生成的库。 例如, 用于 .NET 的被动扩展(Rx) 由一组扩展方法和 LINQ 标准序列运算符组成,以支持异步编程。
何时考虑替代项
IObservable<T>
/
IObserver<T> 接口非常适合基于推送的通知方案,但.NET提供了可能更适合的其他模式:
- Standard .NET 事件 - 对于单个应用程序中的简单通知方案,events更惯用且更易于实现。
-
IAsyncEnumerable<T>— 对于使用者控制节奏的异步拉取序列,请使用 异步流。 -
System.Threading.Channels— 对于具有反压和异步支持的生成者-使用者模式,请使用 System.Threading.Channels。 -
Reactive Extensions (Rx.NET) — 对于复杂的事件组合、筛选和转换,请使用
System.Reactive包,而不是直接实现IObservable<T>。
.NET中最突出的 IObservable<T> 用法是 DiagnosticListener,这使得框架和库作者能够发出使用者订阅的结构化诊断事件。
应用该模式
以下示例使用观察者设计模式来实现机场行李索赔信息系统。
BaggageInfo 类提供有关到达航班以及可领取每次航班行李的行李传送带的信息。 它显示在以下示例中。
namespace Observables.Example;
public readonly record struct BaggageInfo(
int FlightNumber,
string From,
int Carousel);
Namespace Example
Public Structure BaggageInfo
Implements IEquatable(Of BaggageInfo)
Public ReadOnly Property FlightNumber As Integer
Public ReadOnly Property From As String
Public ReadOnly Property Carousel As Integer
Public Sub New(flightNumber As Integer, from As String, carousel As Integer)
Me.FlightNumber = flightNumber
Me.From = from
Me.Carousel = carousel
End Sub
Public Overloads Function Equals(other As BaggageInfo) As Boolean Implements IEquatable(Of BaggageInfo).Equals
Return FlightNumber = other.FlightNumber AndAlso
From = other.From AndAlso
Carousel = other.Carousel
End Function
Public Overrides Function Equals(obj As Object) As Boolean
If TypeOf obj Is BaggageInfo Then
Return Equals(DirectCast(obj, BaggageInfo))
End If
Return False
End Function
Public Overrides Function GetHashCode() As Integer
Return HashCode.Combine(FlightNumber, From, Carousel)
End Function
Public Shared Operator =(left As BaggageInfo, right As BaggageInfo) As Boolean
Return left.Equals(right)
End Operator
Public Shared Operator <>(left As BaggageInfo, right As BaggageInfo) As Boolean
Return Not left.Equals(right)
End Operator
End Structure
End Namespace
BaggageHandler 类负责接收有关到达航班和行李认领传送带的信息。 在内部,它维护两个集合:
-
_observers:观察更新信息的客户端集合。 -
_flights:航班及其指定行李传送带的集合。
以下示例显示了该类的 BaggageHandler 源代码。
namespace Observables.Example;
public sealed class BaggageHandler : IObservable<BaggageInfo>
{
private readonly Lock _lock = new();
private readonly HashSet<IObserver<BaggageInfo>> _observers = [];
private readonly HashSet<BaggageInfo> _flights = [];
public IDisposable Subscribe(IObserver<BaggageInfo> observer)
{
BaggageInfo[] snapshot;
lock (_lock)
{
// Check whether observer is already registered. If not, add it.
if (!_observers.Add(observer))
{
return new Unsubscriber<BaggageInfo>(_lock, _observers, observer);
}
// Snapshot existing data while holding the lock.
snapshot = [.. _flights];
}
// Provide observer with existing data outside the lock.
foreach (BaggageInfo item in snapshot)
{
observer.OnNext(item);
}
return new Unsubscriber<BaggageInfo>(_lock, _observers, observer);
}
// Called to indicate all baggage is now unloaded.
public void BaggageStatus(int flightNumber) =>
BaggageStatus(flightNumber, string.Empty, 0);
public void BaggageStatus(int flightNumber, string from, int carousel)
{
var info = new BaggageInfo(flightNumber, from, carousel);
IObserver<BaggageInfo>[] snapshot;
// Carousel is assigned, so add new info object to list.
if (carousel > 0)
{
lock (_lock)
{
if (!_flights.Add(info))
{
return;
}
snapshot = [.. _observers];
}
foreach (IObserver<BaggageInfo> observer in snapshot)
{
observer.OnNext(info);
}
}
else if (carousel is 0)
{
// Baggage claim for flight is done.
lock (_lock)
{
if (_flights.RemoveWhere(
flight => flight.FlightNumber == info.FlightNumber) == 0)
{
return;
}
snapshot = [.. _observers];
}
foreach (IObserver<BaggageInfo> observer in snapshot)
{
observer.OnNext(info);
}
}
}
public void LastBaggageClaimed()
{
IObserver<BaggageInfo>[] snapshot;
lock (_lock)
{
snapshot = [.. _observers];
_observers.Clear();
}
foreach (IObserver<BaggageInfo> observer in snapshot)
{
observer.OnCompleted();
}
}
}
Namespace Example
Public NotInheritable Class BaggageHandler
Implements IObservable(Of BaggageInfo)
Private ReadOnly _lock As New Object()
Private ReadOnly _observers As New HashSet(Of IObserver(Of BaggageInfo))()
Private ReadOnly _flights As New HashSet(Of BaggageInfo)()
Public Function Subscribe(observer As IObserver(Of BaggageInfo)) As IDisposable Implements IObservable(Of BaggageInfo).Subscribe
Dim snapshot As BaggageInfo()
SyncLock _lock
' Check whether observer is already registered. If not, add it.
If Not _observers.Add(observer) Then
Return New Unsubscriber(Of BaggageInfo)(_lock, _observers, observer)
End If
' Snapshot existing data while holding the lock.
snapshot = _flights.ToArray()
End SyncLock
' Provide observer with existing data outside the lock.
For Each item As BaggageInfo In snapshot
observer.OnNext(item)
Next
Return New Unsubscriber(Of BaggageInfo)(_lock, _observers, observer)
End Function
' Called to indicate all baggage is now unloaded.
Public Sub BaggageStatus(flightNumber As Integer)
BaggageStatus(flightNumber, String.Empty, 0)
End Sub
Public Sub BaggageStatus(flightNumber As Integer, from As String, carousel As Integer)
Dim info As New BaggageInfo(flightNumber, from, carousel)
Dim snapshot As IObserver(Of BaggageInfo)()
' Carousel is assigned, so add new info object to list.
If carousel > 0 Then
SyncLock _lock
If Not _flights.Add(info) Then
Return
End If
snapshot = _observers.ToArray()
End SyncLock
For Each observer As IObserver(Of BaggageInfo) In snapshot
observer.OnNext(info)
Next
ElseIf carousel = 0 Then
' Baggage claim for flight is done.
SyncLock _lock
If _flights.RemoveWhere(
Function(flight) flight.FlightNumber = info.FlightNumber) = 0 Then
Return
End If
snapshot = _observers.ToArray()
End SyncLock
For Each observer As IObserver(Of BaggageInfo) In snapshot
observer.OnNext(info)
Next
End If
End Sub
Public Sub LastBaggageClaimed()
Dim snapshot As IObserver(Of BaggageInfo)()
SyncLock _lock
snapshot = _observers.ToArray()
_observers.Clear()
End SyncLock
For Each observer As IObserver(Of BaggageInfo) In snapshot
observer.OnCompleted()
Next
End Sub
End Class
End Namespace
希望接收更新信息的客户端调用 BaggageHandler.Subscribe 该方法。 如果客户端以前未订阅通知,则将对客户端的 IObserver<T> 实现的引用添加到 _observers 集合中。
可以调用重载 BaggageHandler.BaggageStatus 方法,以指示航班中的行李正在卸载或不再卸载。 在第一种情况下,该方法传递航班号、航班起飞机场和正在卸载行李的传送带。 第二种情况下,该方法仅传递一个航班号。 对于正在卸载的行李,该方法检查传递到方法的 BaggageInfo 信息是否存在于 _flights 集合中。 如果不存在,该方法将添加信息,并调用每个观察者的 OnNext 方法。 对于行李不再卸载的航班,该方法会检查该航班的信息是否存储在_flights集合中。 如果是这样,该方法将调用每个观察者的OnNext 方法,并从BaggageInfo 集合中删除_flights 对象。
当当天的最后一班航班降落并处理行李时,将调用该方法 BaggageHandler.LastBaggageClaimed 。 此方法调用每个观察者 OnCompleted 的方法以指示所有通知都已完成,然后清除 _observers 集合。
提供程序 Subscribe 的方法返回一个 IDisposable 实现,使观察者能够在调用 OnCompleted 方法之前停止接收通知。 此 Unsubscriber 类的源代码在以下示例中显示。 当类在方法中 BaggageHandler.Subscribe 实例化时,它将传递对 _lock 对象、 _observers 集合的引用以及对添加到集合中的观察者的引用。 这些引用将分配给局部变量。 当调用该对象的 Dispose 方法时,它会在锁保护下将观察者从 _observers 集合中移除。
namespace Observables.Example;
internal sealed class Unsubscriber<T> : IDisposable
{
private readonly Lock _lock;
private readonly ISet<IObserver<T>> _observers;
private readonly IObserver<T> _observer;
internal Unsubscriber(
Lock @lock,
ISet<IObserver<T>> observers,
IObserver<T> observer) => (_lock, _observers, _observer) = (@lock, observers, observer);
public void Dispose()
{
lock (_lock)
{
_observers.Remove(_observer);
}
}
}
Namespace Example
Friend NotInheritable Class Unsubscriber(Of T)
Implements IDisposable
Private ReadOnly _lock As Object
Private ReadOnly _observers As ISet(Of IObserver(Of T))
Private ReadOnly _observer As IObserver(Of T)
Friend Sub New(lock As Object, observers As ISet(Of IObserver(Of T)), observer As IObserver(Of T))
_lock = lock
_observers = observers
_observer = observer
End Sub
Public Sub Dispose() Implements IDisposable.Dispose
SyncLock _lock
_observers.Remove(_observer)
End SyncLock
End Sub
End Class
End Namespace
以下示例提供一个名为IObserver<T>ArrivalsMonitor的实现,该实现是显示行李索赔信息的基类。 信息按字母顺序显示,按原始城市的名称显示。 将 ArrivalsMonitor 的方法标记为 overridable(在 Visual Basic 中)或 virtual(在 C# 中),因此它们可在派生类中重写。
namespace Observables.Example;
public class ArrivalsMonitor : IObserver<BaggageInfo>
{
private readonly string _name;
private readonly Lock _lock = new();
private readonly List<string> _flights = [];
private readonly string _format = "{0,-20} {1,5} {2, 3}";
private IDisposable? _cancellation;
public ArrivalsMonitor(string name)
{
ArgumentException.ThrowIfNullOrEmpty(name);
_name = name;
}
public virtual void Subscribe(BaggageHandler provider) =>
_cancellation = provider.Subscribe(this);
public virtual void Unsubscribe()
{
Interlocked.Exchange(ref _cancellation, null)?.Dispose();
lock (_lock)
{
_flights.Clear();
}
}
public virtual void OnCompleted()
{
lock (_lock)
{
_flights.Clear();
}
}
// No implementation needed: Method is not called by the BaggageHandler class.
public virtual void OnError(Exception e)
{
// No implementation.
}
// Update information.
public virtual void OnNext(BaggageInfo info)
{
bool updated = false;
lock (_lock)
{
// Flight has unloaded its baggage; remove from the monitor.
if (info.Carousel is 0)
{
string flightNumber = $"{info.FlightNumber,5}";
for (int index = _flights.Count - 1; index >= 0; index--)
{
string flightInfo = _flights[index];
if (flightInfo.Substring(21, 5).Equals(flightNumber))
{
updated = true;
_flights.RemoveAt(index);
}
}
}
else
{
// Add flight if it doesn't exist in the collection.
string flightInfo = string.Format(_format, info.From, info.FlightNumber, info.Carousel);
if (_flights.Contains(flightInfo) is false)
{
_flights.Add(flightInfo);
updated = true;
}
}
if (updated)
{
_flights.Sort();
Console.WriteLine($"Arrivals information from {_name}");
foreach (string flightInfo in _flights)
{
Console.WriteLine(flightInfo);
}
Console.WriteLine();
}
}
}
}
Imports System.Threading
Namespace Example
Public Class ArrivalsMonitor
Implements IObserver(Of BaggageInfo)
Private ReadOnly _name As String
Private ReadOnly _lock As New Object()
Private ReadOnly _flights As New List(Of String)()
Private ReadOnly _format As String = "{0,-20} {1,5} {2, 3}"
Private _cancellation As IDisposable
Public Sub New(name As String)
If String.IsNullOrEmpty(name) Then
Throw New ArgumentException("Value cannot be null or empty.", NameOf(name))
End If
_name = name
End Sub
Public Overridable Sub Subscribe(provider As BaggageHandler)
_cancellation = provider.Subscribe(Me)
End Sub
Public Overridable Sub Unsubscribe()
Dim previous = Interlocked.Exchange(_cancellation, Nothing)
previous?.Dispose()
SyncLock _lock
_flights.Clear()
End SyncLock
End Sub
Public Overridable Sub OnCompleted() Implements IObserver(Of BaggageInfo).OnCompleted
SyncLock _lock
_flights.Clear()
End SyncLock
End Sub
' No implementation needed: Method is not called by the BaggageHandler class.
Public Overridable Sub OnError([error] As Exception) Implements IObserver(Of BaggageInfo).OnError
' No implementation.
End Sub
' Update information.
Public Overridable Sub OnNext(info As BaggageInfo) Implements IObserver(Of BaggageInfo).OnNext
Dim updated As Boolean = False
SyncLock _lock
' Flight has unloaded its baggage; remove from the monitor.
If info.Carousel = 0 Then
Dim flightNumber As String = String.Format("{0,5}", info.FlightNumber)
For index As Integer = _flights.Count - 1 To 0 Step -1
Dim flightInfo As String = _flights(index)
If flightInfo.Substring(21, 5).Equals(flightNumber) Then
updated = True
_flights.RemoveAt(index)
End If
Next
Else
' Add flight if it doesn't exist in the collection.
Dim flightInfo As String = String.Format(_format, info.From, info.FlightNumber, info.Carousel)
If Not _flights.Contains(flightInfo) Then
_flights.Add(flightInfo)
updated = True
End If
End If
If updated Then
_flights.Sort()
Console.WriteLine($"Arrivals information from {_name}")
For Each flightInfo As String In _flights
Console.WriteLine(flightInfo)
Next
Console.WriteLine()
End If
End SyncLock
End Sub
End Class
End Namespace
该 ArrivalsMonitor 类包括 Subscribe 和 Unsubscribe 方法。
Subscribe 方法使类能够将通过调用 IDisposable 返回的 Subscribe 实现保存到私有变量中。 该方法 Unsubscribe 允许类通过调用提供程序的 Dispose 实现取消订阅通知。
ArrivalsMonitor还提供了OnNext、OnError和OnCompleted方法的实现。 只有OnNext 实现包含大量代码。 该方法与一个专用的、已排序的通用 List<T> 对象配合使用,该对象维护有关抵达航班的始发机场的信息以及其行李可用的传送带。
BaggageHandler如果类报告新的航班到达,则OnNext方法实现会将有关航班的信息添加到列表中。 如果 BaggageHandler 类报告已卸载该航班的行李,则 OnNext 方法从列表中移除该航班。 每当进行更改时,列表将排序并显示到控制台。
以下示例包含应用程序入口点,在此实例化 BaggageHandler 类并创建两个 ArrivalsMonitor 类的实例,使用 BaggageHandler.BaggageStatus 方法添加和删除关于到达航班的信息。 在每个情况下,观察者都会收到更新并正确显示行李索赔信息。
using Observables.Example;
BaggageHandler provider = new();
ArrivalsMonitor observer1 = new("BaggageClaimMonitor1");
ArrivalsMonitor observer2 = new("SecurityExit");
provider.BaggageStatus(712, "Detroit", 3);
observer1.Subscribe(provider);
provider.BaggageStatus(712, "Kalamazoo", 3);
provider.BaggageStatus(400, "New York-Kennedy", 1);
provider.BaggageStatus(712, "Detroit", 3);
observer2.Subscribe(provider);
provider.BaggageStatus(511, "San Francisco", 2);
provider.BaggageStatus(712);
observer2.Unsubscribe();
provider.BaggageStatus(400);
provider.LastBaggageClaimed();
Imports Observables.Example
Imports System.Threading
Module Program
Sub Main(args As String())
Dim provider As New BaggageHandler()
Dim observer1 As New ArrivalsMonitor("BaggageClaimMonitor1")
Dim observer2 As New ArrivalsMonitor("SecurityExit")
provider.BaggageStatus(712, "Detroit", 3)
observer1.Subscribe(provider)
provider.BaggageStatus(712, "Kalamazoo", 3)
provider.BaggageStatus(400, "New York-Kennedy", 1)
provider.BaggageStatus(712, "Detroit", 3)
observer2.Subscribe(provider)
provider.BaggageStatus(511, "San Francisco", 2)
provider.BaggageStatus(712)
observer2.Unsubscribe()
provider.BaggageStatus(400)
provider.LastBaggageClaimed()
End Sub
End Module