Postup implementace zprostředkovatele

Vzor návrhu pozorovatele vyžaduje rozdělení mezi poskytovatelem, který monitoruje data a odesílá oznámení, a jeden nebo více pozorovatelů, kteří od poskytovatele přijímají oznámení (zpětné volání). Tento článek ukazuje, jak vytvořit zprostředkovatele. Informace o vytvoření pozorovatele naleznete v tématu Postup implementace pozorovatele.

Definování datového typu

Definujte data, která poskytovatel odesílá pozorovatelům. Přestože zprostředkovatel a data, která odesílá pozorovatelům, mohou být jedním typem, obvykle každý z nich představuje jiný typ. Například v aplikaci pro monitorování teploty struktura Temperature definuje data, která třída TemperatureMonitor (definovaná v následující části) monitoruje a k jejichž odběru se pozorovatelé přihlašují.

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

Vytvoření zprostředkovatele

Zprostředkovatel dat je typ, který implementuje System.IObservable<T> rozhraní. Argument obecného typu poskytovatele je typ, který odesílá pozorovatelům.

  1. Definujte třídu zprostředkovatele. Následující příklad definuje TemperatureMonitor třídu, což je konstruovaná System.IObservable<T> implementace s obecným typem argumentu 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. Přidejte pole pro ukládání referencí na pozorovatele.

    Poskytovatel musí sledovat každého registrovaného pozorovatele, aby mohl odesílat oznámení později. Obvykle použijte objekt kolekce, jako je obecný List<T> objekt. Následující příklad definuje privátní objekt List<T>, který je vytvořen v konstruktoru třídy TemperatureMonitor.

    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. IDisposable Definujte implementaci pro zrušení odběru.

    Poskytovatel tuto implementaci vrátí odběratelům, aby mohli kdykoli ukončit příjem oznámení. Následující příklad definuje vnořenou třídu Unsubscriber, která při vytvoření instance přebírá odkaz na kolekci odběratelů a na odběratele. Třída Unsubscriber umožňuje odběrateli volat implementaci IDisposable.Dispose objektu, aby se odhlásil z kolekce odběratelů.

    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. Implementujte metodu IObservable<T>.Subscribe .

    Metoda obdrží odkaz na System.IObserver<T> rozhraní. Uložte tuto referenci do kolekce pozorovatelů z předchozího kroku a potom vraťte implementaci pro odhlášení odběru IDisposable. Následující příklad ukazuje implementaci Subscribe ve TemperatureMonitor třídě.

    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. Implementujte logiku upozorňování voláním metod pozorovatelů IObserver<T>.OnNext, IObserver<T>.OnError a IObserver<T>.OnCompleted.

    V některých případech nemusí poskytovatel volat OnError , když dojde k chybě. Následující GetTemperature metoda simuluje monitorování, které čte data o teplotě každých pět sekund, a upozorní pozorovatele, pokud se teplota od předchozího čtení změnila alespoň o 0,1 stupňů. Pokud zařízení neohlásí teplotu (tj. pokud je její hodnota null), poskytovatel upozorní pozorovatele, že přenos je dokončen voláním metody každého pozorovatele OnCompleted a vymaže List<T> kolekci. V tomto příkladu poskytovatel nikdy nevolá 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
    

Příklad

Následující příklad obsahuje úplný zdrojový kód pro IObservable<T> implementaci aplikace pro monitorování teploty. Zahrnuje Temperature strukturu, což je data, která poskytovatel odesílá pozorovatelům, a TemperatureMonitor třídu, což je IObservable<T> implementace.

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