言語

プロバイダーを実装する方法

オブザーバーの設計パターンでは、データを監視して通知を送信するプロバイダーと、プロバイダーから通知 (コールバック) を受信する 1 つ以上のオブザーバーとの間で除算が必要です。 この記事では、プロバイダーを作成する方法について説明します。 オブザーバーの作成の詳細については、「オブザーバー を実装する方法」を参照してください。

データ型を定義する

プロバイダーがオブザーバーに送信するデータを定義します。 プロバイダーとオブザーバーに送信するデータは 1 つの型にできますが、通常は異なる型がそれぞれを表します。 たとえば、温度監視アプリケーションでは、 Temperature 構造体は、 TemperatureMonitor クラス (次のセクションで定義) が監視し、オブザーバーがサブスクライブするデータを定義します。

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

プロバイダーの作成

データ プロバイダーは、 System.IObservable<T> インターフェイスを実装する型です。 プロバイダーのジェネリック型引数は、オブザーバーに送信される型です。

  1. プロバイダー クラスを定義します。 次の例では、TemperatureMonitor クラスを定義します。これは、System.IObservable<T>のジェネリック型引数を使用して構築された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. オブザーバー参照を格納するフィールドを追加します。

    プロバイダーは、後で通知を送信できるように、登録された各オブザーバーを追跡する必要があります。 通常は、ジェネリック List<T> オブジェクトなどのコレクション オブジェクトを使用します。 次の例では、List<T> クラス コンストラクターでインスタンス化されたプライベート 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 実装を定義します。

    プロバイダーは、いつでも通知の受信を停止できるように、サブスクライバーにこの実装を返します。 次の例では、インスタンス化されたときにサブスクライバー コレクションとサブスクライバーへの参照を受け取る入れ子になった Unsubscriber クラスを定義します。 Unsubscriber クラスを使用すると、サブスクライバーはオブジェクトのIDisposable.Dispose実装を呼び出して、サブスクライバー コレクションから自身を削除できます。

    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. IObservable<T>.Subscribe メソッドを実装します。

    メソッドは、 System.IObserver<T> インターフェイスへの参照を受け取ります。 その参照を前の手順のオブザーバー コレクションに格納し、その後 IDisposable の登録解除の実装を返します。 次の例は、Subscribe クラスのTemperatureMonitor実装を示しています。

    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. オブザーバーの IObserver<T>.OnNextIObserver<T>.OnErrorIObserver<T>.OnCompleted メソッドを呼び出して、通知ロジックを実装します。

    場合によっては、エラーが発生したときにプロバイダーが OnError を呼び出さない場合があります。 次の GetTemperature メソッドは、5 秒ごとに温度データを読み取り、前の読み取り以降に温度が少なくとも .1 度変化した場合にオブザーバーに通知するモニターをシミュレートします。 デバイスが温度を報告しない場合 (つまり、値が null の場合)、プロバイダーは、各オブザーバーの OnCompleted メソッドを呼び出して送信が完了したことをオブザーバーに通知し、 List<T> コレクションをクリアします。 この例では、プロバイダーは 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
    

次の例には、温度監視アプリケーションの IObservable<T> 実装の完全なソース コードが含まれています。 これには、プロバイダーがオブザーバーに送信するデータであるTemperature構造体と、TemperatureMonitor実装である IObservable<T> クラスが含まれます。

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