言語
オブザーバーの設計パターンでは、データを監視して通知を送信するプロバイダーと、プロバイダーから通知 (コールバック) を受信する 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> インターフェイスを実装する型です。 プロバイダーのジェネリック型引数は、オブザーバーに送信される型です。
プロバイダー クラスを定義します。 次の例では、
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)オブザーバー参照を格納するフィールドを追加します。
プロバイダーは、後で通知を送信できるように、登録された各オブザーバーを追跡する必要があります。 通常は、ジェネリック 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()サブスクライブを解除するための 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 ClassIObservable<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オブザーバーの IObserver<T>.OnNext、 IObserver<T>.OnError、 IObserver<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
関連するコンテンツ
.NET