Sprache

So implementieren Sie einen Anbieter

Das Entwurfsmuster des Beobachters erfordert eine Trennung zwischen einem Anbieter, der Daten überwacht und Benachrichtigungen sendet, und einem oder mehreren Beobachtern, die Benachrichtigungen (Rückrufe) vom Anbieter empfangen. In diesem Artikel wird gezeigt, wie Sie einen Anbieter erstellen. Informationen zum Erstellen eines Beobachters finden Sie unter Implementieren eines Beobachters.

Definieren des Datentyps

Definieren Sie die Daten, die der Anbieter an Beobachter sendet. Obwohl der Anbieter und die Daten, die er an Beobachter sendet, ein einzelner Typ sein können, stellt ein anderer Typ normalerweise jeden dar. Zum Beispiel definiert in einer Temperaturüberwachungsanwendung die Temperature-Struktur die Daten, die die TemperatureMonitor-Klasse (die im nächsten Abschnitt definiert wird) überwacht und die von Beobachtern abonniert werden.

namespace TemperatureSample;

public readonly record struct Temperature(decimal Degrees, DateTime Date);
Namespace Global.TemperatureSample

    Public Structure Temperature
        Public ReadOnly Property Degrees As Decimal
        Public ReadOnly Property [Date] As Date

        Public Sub New(degrees As Decimal, [date] As Date)
            Me.Degrees = degrees
            Me.Date = [date]
        End Sub
    End Structure

End Namespace

Erstellen eines Anbieters

Der Datenanbieter ist ein Typ, der die System.IObservable<T> Schnittstelle implementiert. Das generische Typargument des Anbieters ist der Typ, den er an Beobachter sendet.

  1. Definieren Sie die Anbieterklasse. Im folgenden Beispiel wird eine TemperatureMonitor Klasse definiert, bei der es sich um eine konstruierte System.IObservable<T> Implementierung mit einem generischen Typargument handelt Temperature.

    namespace TemperatureSample;
    
    public sealed class TemperatureMonitor : IObservable<Temperature>
    {
    
    Imports System.Threading
    Imports System.Threading.Tasks
    
    Namespace Global.TemperatureSample
    
        Public NotInheritable Class TemperatureMonitor
            Implements IObservable(Of Temperature)
    
  2. Fügen Sie ein Feld hinzu, um Beobachterverweise zu speichern.

    Der Anbieter muss jeden registrierten Beobachter nachverfolgen, damit er später Benachrichtigungen senden kann. Verwenden Sie in der Regel ein Auflistungsobjekt wie ein generisches List<T> Objekt. Im folgenden Beispiel wird ein privates List<T> Objekt definiert, das im Konstruktor der Klasse TemperatureMonitor instanziiert wird.

    namespace TemperatureSample;
    
    public sealed class TemperatureMonitor : IObservable<Temperature>
    {
        private readonly List<IObserver<Temperature>> _observers = [];
        private readonly Lock _sync = new();
    
    Imports System.Threading
    Imports System.Threading.Tasks
    
    Namespace Global.TemperatureSample
    
        Public NotInheritable Class TemperatureMonitor
            Implements IObservable(Of Temperature)
    
            Private ReadOnly _observers As New List(Of IObserver(Of Temperature))()
            Private ReadOnly _sync As New Object()
    
  3. Definieren Sie eine IDisposable Implementierung für das Abmelden.

    Der Anbieter gibt diese Implementierung an Abonnenten zurück, damit sie jederzeit den Empfang von Benachrichtigungen beenden können. Im folgenden Beispiel wird eine geschachtelte Unsubscriber Klasse definiert, die beim Instanziieren einen Verweis auf die Abonnentensammlung und den Abonnenten empfängt. Die Unsubscriber Klasse ermöglicht es dem Abonnenten, die Implementierung des IDisposable.Dispose Objekts aufzurufen, um sich selbst aus der Abonnentensammlung zu entfernen.

    private sealed class Unsubscriber(
        List<IObserver<Temperature>> observers,
        IObserver<Temperature> observer,
        Lock sync) : IDisposable
    {
        public void Dispose()
        {
            lock (sync)
            {
                observers.Remove(observer);
            }
        }
    }
    
    Private NotInheritable Class Unsubscriber
        Implements IDisposable
    
        Private ReadOnly _observers As List(Of IObserver(Of Temperature))
        Private ReadOnly _observer As IObserver(Of Temperature)
        Private ReadOnly _sync As Object
    
        Public Sub New(observers As List(Of IObserver(Of Temperature)),
                       observer As IObserver(Of Temperature),
                       sync As Object)
            _observers = observers
            _observer = observer
            _sync = sync
        End Sub
    
        Public Sub Dispose() Implements IDisposable.Dispose
            SyncLock _sync
                _observers.Remove(_observer)
            End SyncLock
        End Sub
    End Class
    
  4. Implementieren Sie die IObservable<T>.Subscribe-Methode.

    Die Methode empfängt einen Verweis auf die System.IObserver<T> Schnittstelle. Speichern Sie diesen Verweis in der Beobachtersammlung aus dem vorherigen Schritt, und geben Sie dann die IDisposable Abbestellerimplementierung zurück. Das folgende Beispiel zeigt die Subscribe Implementierung in der TemperatureMonitor Klasse.

    public IDisposable Subscribe(IObserver<Temperature> observer)
    {
        ArgumentNullException.ThrowIfNull(observer);
    
        lock (_sync)
        {
            if (!_observers.Contains(observer))
                _observers.Add(observer);
        }
    
        return new Unsubscriber(_observers, observer, _sync);
    }
    
    Public Function Subscribe(observer As IObserver(Of Temperature)) As IDisposable _
        Implements IObservable(Of Temperature).Subscribe
    
        ArgumentNullException.ThrowIfNull(observer)
    
        SyncLock _sync
            If Not _observers.Contains(observer) Then
                _observers.Add(observer)
            End If
        End SyncLock
    
        Return New Unsubscriber(_observers, observer, _sync)
    End Function
    
  5. Implementieren Sie die Benachrichtigungslogik, indem Sie die Methoden IObserver<T>.OnNext, IObserver<T>.OnError und IObserver<T>.OnCompleted der Beobachter aufrufen.

    In einigen Fällen ruft ein Anbieter OnError möglicherweise nicht auf, wenn ein Fehler auftritt. Die folgende GetTemperature Methode simuliert einen Monitor, der Temperaturdaten alle fünf Sekunden liest, und benachrichtigt Beobachter, wenn sich die Temperatur seit dem vorherigen Lesen um mindestens 0,1 Grad geändert hat. Wenn das Gerät keine Temperatur meldet (d. h., wenn der Wert null ist), benachrichtigt der Anbieter Beobachter darüber, dass die Übertragung abgeschlossen ist, indem die Methode der einzelnen Beobachter OnCompleted aufgerufen und die List<T> Sammlung gelöscht wird. In diesem Beispiel ruft der Anbieter nie auf OnError.

    public async Task GetTemperatureAsync(CancellationToken cancellationToken = default)
    {
        // Sample data that mimics a temperature device. A null value signals the end of transmission.
        decimal?[] temps =
        [
            14.6m, 14.65m, 14.7m, 14.9m, 14.9m, 15.2m,
            15.25m, 15.2m, 15.4m, 15.45m, null
        ];
    
        decimal? previous = null;
    
        foreach (decimal? temp in temps)
        {
            await Task.Delay(TimeSpan.FromSeconds(2.5), cancellationToken);
    
            if (temp is decimal value)
            {
                // Notify only after at least a 0.1° change.
                if (previous is null || Math.Abs(value - previous.Value) >= 0.1m)
                {
                    NotifyAll(new Temperature(value, DateTime.Now));
                    previous = value;
                }
            }
            else
            {
                CompleteAll();
                break;
            }
        }
    }
    
    private void NotifyAll(Temperature data)
    {
        IObserver<Temperature>[] snapshot;
        lock (_sync)
        {
            snapshot = [.. _observers];
        }
    
        foreach (IObserver<Temperature> observer in snapshot)
            observer.OnNext(data);
    }
    
    private void CompleteAll()
    {
        IObserver<Temperature>[] snapshot;
        lock (_sync)
        {
            snapshot = [.. _observers];
            _observers.Clear();
        }
    
        foreach (IObserver<Temperature> observer in snapshot)
            observer.OnCompleted();
    }
    
    Public Async Function GetTemperatureAsync(Optional cancellationToken As CancellationToken = Nothing) As Task
        ' Sample data that mimics a temperature device. A Nothing value signals the end of transmission.
        Dim temps As Decimal?() = {
            14.6D, 14.65D, 14.7D, 14.9D, 14.9D, 15.2D,
            15.25D, 15.2D, 15.4D, 15.45D, Nothing
        }
    
        Dim previous As Decimal? = Nothing
    
        For Each temp As Decimal? In temps
            Await Task.Delay(TimeSpan.FromSeconds(2.5), cancellationToken)
    
            If temp.HasValue Then
                ' Notify only after at least a 0.1° change.
                If Not previous.HasValue OrElse Math.Abs(temp.Value - previous.Value) >= 0.1D Then
                    NotifyAll(New Temperature(temp.Value, Date.Now))
                    previous = temp
                End If
            Else
                CompleteAll()
                Exit For
            End If
        Next
    End Function
    
    Private Sub NotifyAll(data As Temperature)
        Dim snapshot As IObserver(Of Temperature)()
        SyncLock _sync
            snapshot = _observers.ToArray()
        End SyncLock
    
        For Each observer In snapshot
            observer.OnNext(data)
        Next
    End Sub
    
    Private Sub CompleteAll()
        Dim snapshot As IObserver(Of Temperature)()
        SyncLock _sync
            snapshot = _observers.ToArray()
            _observers.Clear()
        End SyncLock
    
        For Each observer In snapshot
            observer.OnCompleted()
        Next
    End Sub
    

Example

Das folgende Beispiel enthält den vollständigen Quellcode für eine IObservable<T> Implementierung für eine Temperaturüberwachungsanwendung. Sie enthält die Temperature Struktur, bei der es sich um die Daten handelt, die der Anbieter an Beobachter sendet, und die TemperatureMonitor Klasse, die die IObservable<T> Implementierung ist.

namespace TemperatureSample;

public sealed class TemperatureMonitor : IObservable<Temperature>
{
    private readonly List<IObserver<Temperature>> _observers = [];
    private readonly Lock _sync = new();

    private sealed class Unsubscriber(
        List<IObserver<Temperature>> observers,
        IObserver<Temperature> observer,
        Lock sync) : IDisposable
    {
        public void Dispose()
        {
            lock (sync)
            {
                observers.Remove(observer);
            }
        }
    }

    public IDisposable Subscribe(IObserver<Temperature> observer)
    {
        ArgumentNullException.ThrowIfNull(observer);

        lock (_sync)
        {
            if (!_observers.Contains(observer))
                _observers.Add(observer);
        }

        return new Unsubscriber(_observers, observer, _sync);
    }

    public async Task GetTemperatureAsync(CancellationToken cancellationToken = default)
    {
        // Sample data that mimics a temperature device. A null value signals the end of transmission.
        decimal?[] temps =
        [
            14.6m, 14.65m, 14.7m, 14.9m, 14.9m, 15.2m,
            15.25m, 15.2m, 15.4m, 15.45m, null
        ];

        decimal? previous = null;

        foreach (decimal? temp in temps)
        {
            await Task.Delay(TimeSpan.FromSeconds(2.5), cancellationToken);

            if (temp is decimal value)
            {
                // Notify only after at least a 0.1° change.
                if (previous is null || Math.Abs(value - previous.Value) >= 0.1m)
                {
                    NotifyAll(new Temperature(value, DateTime.Now));
                    previous = value;
                }
            }
            else
            {
                CompleteAll();
                break;
            }
        }
    }

    private void NotifyAll(Temperature data)
    {
        IObserver<Temperature>[] snapshot;
        lock (_sync)
        {
            snapshot = [.. _observers];
        }

        foreach (IObserver<Temperature> observer in snapshot)
            observer.OnNext(data);
    }

    private void CompleteAll()
    {
        IObserver<Temperature>[] snapshot;
        lock (_sync)
        {
            snapshot = [.. _observers];
            _observers.Clear();
        }

        foreach (IObserver<Temperature> observer in snapshot)
            observer.OnCompleted();
    }
}
Imports System.Threading
Imports System.Threading.Tasks

Namespace Global.TemperatureSample

    Public NotInheritable Class TemperatureMonitor
        Implements IObservable(Of Temperature)

        Private ReadOnly _observers As New List(Of IObserver(Of Temperature))()
        Private ReadOnly _sync As New Object()

        Private NotInheritable Class Unsubscriber
            Implements IDisposable

            Private ReadOnly _observers As List(Of IObserver(Of Temperature))
            Private ReadOnly _observer As IObserver(Of Temperature)
            Private ReadOnly _sync As Object

            Public Sub New(observers As List(Of IObserver(Of Temperature)),
                           observer As IObserver(Of Temperature),
                           sync As Object)
                _observers = observers
                _observer = observer
                _sync = sync
            End Sub

            Public Sub Dispose() Implements IDisposable.Dispose
                SyncLock _sync
                    _observers.Remove(_observer)
                End SyncLock
            End Sub
        End Class

        Public Function Subscribe(observer As IObserver(Of Temperature)) As IDisposable _
            Implements IObservable(Of Temperature).Subscribe

            ArgumentNullException.ThrowIfNull(observer)

            SyncLock _sync
                If Not _observers.Contains(observer) Then
                    _observers.Add(observer)
                End If
            End SyncLock

            Return New Unsubscriber(_observers, observer, _sync)
        End Function

        Public Async Function GetTemperatureAsync(Optional cancellationToken As CancellationToken = Nothing) As Task
            ' Sample data that mimics a temperature device. A Nothing value signals the end of transmission.
            Dim temps As Decimal?() = {
                14.6D, 14.65D, 14.7D, 14.9D, 14.9D, 15.2D,
                15.25D, 15.2D, 15.4D, 15.45D, Nothing
            }

            Dim previous As Decimal? = Nothing

            For Each temp As Decimal? In temps
                Await Task.Delay(TimeSpan.FromSeconds(2.5), cancellationToken)

                If temp.HasValue Then
                    ' Notify only after at least a 0.1° change.
                    If Not previous.HasValue OrElse Math.Abs(temp.Value - previous.Value) >= 0.1D Then
                        NotifyAll(New Temperature(temp.Value, Date.Now))
                        previous = temp
                    End If
                Else
                    CompleteAll()
                    Exit For
                End If
            Next
        End Function

        Private Sub NotifyAll(data As Temperature)
            Dim snapshot As IObserver(Of Temperature)()
            SyncLock _sync
                snapshot = _observers.ToArray()
            End SyncLock

            For Each observer In snapshot
                observer.OnNext(data)
            Next
        End Sub

        Private Sub CompleteAll()
            Dim snapshot As IObserver(Of Temperature)()
            SyncLock _sync
                snapshot = _observers.ToArray()
                _observers.Clear()
            End SyncLock

            For Each observer In snapshot
                observer.OnCompleted()
            Next
        End Sub
    End Class

End Namespace