はじめに
誤りについては優しく指摘していただけますと幸いです(まだまだ理解できていないので誤りだらけかと思います)。
Reactive Extensions (Rx) とは何か
一言でいうと、状態変化に応じた処理が書きやすいことが利点。
対象を観測し、その状態が変化すると観測側に通知が行くシステム。
いつ何回発生するか不明なイベントがあり、
それについて処理を行うシチュエーションを考える。
LINQのように処理する側 = データを要求する側に主導権があるとスレッドをブロックしてしまう。 これは同期的に値を返すことを前提としているために引き起こされる。(Pull型)
一方Rxではイベント側 = データを提供する側に主導権があるため、スレッドをブロックせずに済む。 これは、値が来たら提供側から届けるという契約になっているため。(Push型)
static void Main(string[] args)
{
Console.WriteLine("開始");
// この後3秒間はスレッドがブロックされる
foreach (var x in Events())
{
Console.WriteLine(x);
}
Console.WriteLine("終了");
Console.ReadKey();
}
static IEnumerable<int> Events()
{
// イベントを3秒待つ想定
Thread.Sleep(3000);
yield return 1;
}
static void Main()
{
Console.WriteLine("開始");
// 3秒後に「1」というデータを通知し、表示する
Observable.Timer(TimeSpan.FromSeconds(3))
.Subscribe(x => Console.WriteLine("1"));
// スレッドがブロックされないため、3秒待たずに「終了」と表示される。
Console.WriteLine("終了");
Console.ReadKey();
}
実行例
似たようなことがAsync/Awaitでも可能だが、あちらは連続して起こる処理には向かない。
「状態変化」はいつ起きるか分からない上に何度も発生しうる。
そんな時にはRxが威力を発揮する。
Rxは通知を便利に扱うライブラリ。
なのでいたずらに使うと、値が可変になってしまい クラスの安定性が低下する。
利用する側の実装難易度が上がってしまうというリスクも持っている。
使い分けまとめ
| LINQ | Async/Await | Rx |
|---|---|---|
| DBから取得して 加工 |
WebAPIを1回呼び出す | UIのイベント処理 |
| メモリ上の 配列の集計 |
順序だった非同期処理 | リアルタイム通知の処理 |
| CSV/JSONの 読み込み |
並行して処理し、 全て終了するのを待つ |
WebSocketのプッシュ通知 |
用語のイメージ
| 用語 | イメージ |
|---|---|
| ストリーム | 川の流れ |
| Observable | 川の設計・水源 |
| Observer | 川下で水を受け取る人 |
| subscribe | 水門の開放 |
| subject | 誰でも水を注ぎ・汲める池 |
| 実データ | 水 |
Where/Select/Subscribe
基本的なRxのメソッドの具体例。
static void Main()
{
using (NumberWriter numberWriter = new NumberWriter())
using (IDisposable subscriber = numberWriter.Writer
.Where(n => n >= 5)
.Select(n => n + 100)
.Subscribe(
onNext: number => Console.WriteLine($"受信: {number}"),
onCompleted: () => Console.WriteLine("完了")
)
)
{
numberWriter.CallAndAdd(1);
numberWriter.CallAndAdd(3);
numberWriter.CallAndAdd(5);
numberWriter.CallAndAdd(7);
numberWriter.Terminate();
numberWriter.CallAndAdd(9);
Console.ReadKey();
}
}
public class NumberWriter : IDisposable
{
public Subject<int> _writer = new Subject<int>();
public IObservable<int> Writer => _writer;
public void CallAndAdd(int number)
{
_writer.OnNext(number);
}
public void Terminate()
{
_writer.OnCompleted();
}
public void Dispose()
{
_writer.Dispose();
}
}
Where
LINQのWhereと同様に、流れて来るデータをフィルタリングする。
Select
LINQのSelectと同様に、流れて来るデータを加工する。
Subscribe
データを流し始めるという宣言。これが呼ばれた部分からデータが流れるので、これ以前に流されたデータは渡されない。戻り値はIDisposable。
IObservable<T>が呼び出すことができ、 引数にはIObserver<T>が入る。
-
onNextには変化が通知された場合に行う処理が書かれており、
IObservableはIObserver(つまり書かれた処理)にデータを渡し、処理を行う。 -
onCompletedはデータが流れ終わったときに行われる処理が書かれており、
これ以降に流れたデータは通知されない。 -
onErrorというものもあり、こちらはエラーの通知が来た時に行う処理が書かれる。
GroupBy
static void Main()
{
var source = Observable.Range(1, 10);
source.GroupBy(v => v % 2)
.Subscribe(g =>
{
IGroupedObservable<int, int> group = g;
group.Subscribe(value => Console.WriteLine($"グループ: {group.Key}, {value}"));
});
}
GropuByは
- 流れてきたストリームの分割
- 分割されたグループに値を流す
という2つのことを行っていると考えている。
つまり、sourceで1~10までの数値というストリームが作られてGroupByに流れて来る。
- まず1が流れてくると奇数のグループとして
IGroupedObservableを生成する - その後、そのグループには「グループの
Keyと流れてきた値を出力する」という処理を登録して1を渡す。すると、グループのKey(今回は2で割ったあまり) と渡された1が表示される - 次に2が渡された際の処理も同様。3が渡された場合は奇数グループになるが、そのグループは既に存在するため新規の作成は行われない。その工程はスキップされ、そのままストリームが流れていく (データが渡される) 形になる
FromEvent
WinFormsなどの既存ライブラリは通知をeventで提供している。
しかし、WhereやSelectといった演算子はIObservable<T>にしか使えない。
既存のイベントをRxで扱う手段として、FromEventが挙げられる。
これにより、eventをIObservable<T>という「値」として扱えるようになり、
Whereや後述のMergeに渡せるようになる。
また、eventはKeyEventHandler、EventHandler、MouseEventHandlerなどイベントごとに異なるシグネチャを持っているが、
これを統一的に扱えるようになることもメリットとして挙げられる。
以下のサンプルコードではFormが表示され、
Enterキーを押した時にだけコンソールにReturnと表示される。
var form = new Form();
Observable.FromEvent<KeyEventHandler, KeyEventArgs>(
onNext => (sender, eventArgs) => onNext(eventArgs),
handler => form.KeyDown += handler,
handler => form.KeyDown -= handler
).Where(eventArgs => eventArgs.KeyCode == Keys.Enter
).Subscribe(e => Console.WriteLine(e.KeyCode)
);
Application.Run(form);
最初にこれを見せられても良く分からなかったので、
以下にラムダ式を減らして書いたバージョンも記載する。
class Program
{
static Form form = new Form();
static void Main()
{
Observable.FromEvent<KeyEventHandler, KeyEventArgs>(
ConvertToHandler,
AddHandler,
SubtractHandler
)
.Where(eventArgs => eventArgs.KeyCode == Keys.Enter)
.Subscribe(e => Console.WriteLine(e.KeyCode)
);
Application.Run(form);
}
/// <summary>
/// 通知用のActionをEventが要求する形に変換する
/// </summary>
/// <param name="keyEventAction"></param>
/// <returns></returns>
static KeyEventHandler ConvertToHandler(Action<KeyEventArgs> keyEventAction)
{
Console.WriteLine("Convert が呼ばれた");
return new KeyEventHandler(
(object sender, KeyEventArgs eventArgs) => keyEventAction(eventArgs)
);
}
static void AddHandler(KeyEventHandler keyEventHandler)
{
Console.WriteLine("Add が呼ばれた");
form.KeyDown += keyEventHandler;
}
static void SubtractHandler(KeyEventHandler keyEventHandler)
{
Console.WriteLine("Subtract が呼ばれた");
form.KeyDown -= keyEventHandler;
}
}
現時点での理解を記しておくが、読みづらいので後々読みやすい形にする予定。
-
FromEventでAction<KeyEventArgs>がKeyEventHandlerに変換する処理と、formのKeyDownイベントにhandlerを登録する処理が保持される。これらはSubscribeされた瞬間に登録処理が走る -
Action<KeyEventArgs>に渡されるのは、今回の場合Subscribeの中に書かれたIObserverのラムダ式になる -
formが表示されるとKeyを押すというイベントを受け付けるようになる -
Keyを押す(ストリームが流れる)とSubscribeの内部に書かれたIObserverのonNextを呼び、そこにキーの情報を渡そうとする。その際、まずWhereの判定にかかり、Enterか否かが判断される。 - Enterでなければそれ以降にストリームは流されず、Enterであればストリームが流れていく。
Subscribeの中に書かれたIObserverにそれは渡され、onNextの処理 =Subscribeに書いたラムダ式の処理が行われる。最終的にReturnが出力される
Merge
キャンセル処理、購入ボタン、複数センサの異常検知など、
発生源が異なるが受け取った後の処理が同じ場合、
発生源ごとに処理を書く必要が出てくる。
それを回避することができるのがMerge。
static void Main()
{
Subject<Unit> buttonA = new Subject<Unit>();
Subject<Unit> buttonB = new Subject<Unit>();
int totalCount = 0;
buttonA.Merge(buttonB)
.Subscribe(_ => {
totalCount++;
Console.WriteLine($"{totalCount}回クリック");
}
);
ClickWithLog(buttonA, "ButtonA");
ClickWithLog(buttonB, "ButtonB");
ClickWithLog(buttonA, "ButtonA");
ClickWithLog(buttonB, "ButtonB");
ClickWithLog(buttonA, "ButtonA");
}
public static void ClickWithLog(Subject<Unit> button, string name)
{
Console.WriteLine($"Click: {name}");
button.OnNext(Unit.Default);
}
buttonAとbuttonBはどちらも独立したSubjectのままであり続けるが、
Merge内部では新しいIOPbservable<Unit>が作成される。
新しいIObservable<Unit>はSubscribeされた時点でbuttonA・Bの両方を購読する。
Subject は、内部に登録された処理=IObserverのリストを持っている。OnNextが呼ばれると、そのリストの全員のOnNextを順に呼ぶ。
つまり、buttonのOnNextが呼ばれると、
IObservableのOnNextが呼ばれるということになる。
どちらかのOnNextが呼ばれると、
その値がMergeに流れてSubscribe内部のラムダ式が実行される。
CombineLatest
ログイン画面などでユーザー名とパスワードの両方が入力されるまで処理を行いたくないという場合がある。あるいは、温度と湿度の最新値から不快指数を計算するため、両者が揃うまで処理をおこないたくないという場面がある。
そのような場面で使えるのがCombineLatest。
CombineLatesetは各入力の最新値を覚えており、それに基づいて処理を行う。
全入力が1回ずつ値を出すまでは何もしないのが特徴。
static void Main()
{
Subject<string> userNameInput = new Subject<string>();
Subject<string> passwordInput = new Subject<string>();
userNameInput.CombineLatest(passwordInput,
(name, pass) => !string.IsNullOrEmpty(name) && !string.IsNullOrEmpty(pass)
? "ログイン可能" : "ログイン不可")
.Subscribe(v => Console.WriteLine($"Result: {v}"));
InputWithLog(userNameInput, "yamada");
InputWithLog(passwordInput, "");
InputWithLog(passwordInput, "password");
InputWithLog(userNameInput, "");
InputWithLog(userNameInput, "yamada");
}
public static void InputWithLog(Subject<string> input, string value)
{
Console.WriteLine($"Input: {value}");
input.OnNext(value);
}
DistinctUntilChanged と Distinct
直前と同じ値の場合は通知されない。
「変化」のみを伝えることで、無駄な通知を省くことができる。
static void Main()
{
var temperatures = new[] { 20, 20, 21, 21, 21, 20, 22, 22, 23, 20, 20 };
using (TemperatureDisplayer displayer = new TemperatureDisplayer())
using (IDisposable subscriber = displayer.Subject
.DistinctUntilChanged(
).Subscribe(
onNext: value => Console.WriteLine($"気温: {value}"),
onCompleted: () => Console.WriteLine("出力終了"))
)
{
Console.WriteLine("センサデータ送信開始");
foreach(var temp in temperatures)
{
displayer.OnNext(temp);
}
displayer.OnCompleted();
}
}
public class TemperatureDisplayer : IDisposable
{
public Subject<int> _subject = new Subject<int>();
public IObservable<int> Subject => _subject;
public void OnNext(int value)
{
_subject.OnNext(value);
}
public void OnCompleted()
{
_subject.OnCompleted();
}
public void Dispose()
{
_subject.Dispose();
}
}
類似したものにDistinctがあるが、こちらは今までに通知していない値のみを通知する。
今までに1度でも出ている値は通知されない。
挙動が異なるので注意。
下のコード例では気温にのみ基づいて同じ値か否かを比較するため、
data => data.Temperatureと記載している。
static void Main()
{
var sensorDataLine = new[]
{
new SensorData(new DateTime(2026, 1, 1, 9, 0, 0), 20),
new SensorData(new DateTime(2026, 1, 1, 9, 1, 0), 20),
new SensorData(new DateTime(2026, 1, 1, 9, 2, 0), 21),
new SensorData(new DateTime(2026, 1, 1, 9, 3, 0), 21),
new SensorData(new DateTime(2026, 1, 1, 9, 4, 0), 22),
new SensorData(new DateTime(2026, 1, 1, 9, 5, 0), 20),
new SensorData(new DateTime(2026, 1, 1, 9, 6, 0), 23),
};
using(SensorDataDisplayer displayer = new SensorDataDisplayer())
using (IDisposable subscriber = displayer.Subject
.Distinct(data => data.Temperature)
.Subscribe(
onNext: data => Console.WriteLine($"気温: {data.Temperature} (観測時刻: {data.Time})"),
onCompleted: () => Console.WriteLine("出力終了"))
)
{
Console.WriteLine("センサデータ送信開始");
foreach(var temp in sensorDataLine)
{
displayer.OnNext(temp);
}
displayer.OnCompleted();
}
}
public class SensorDataDisplayer : IDisposable
{
public Subject<SensorData> _subject = new Subject<SensorData>();
public IObservable<SensorData> Subject => _subject;
public void OnNext(SensorData data)
{
_subject.OnNext(data);
}
public void OnCompleted()
{
_subject.OnCompleted();
}
public void Dispose()
{
_subject.Dispose();
}
}
public class SensorData
{
public DateTime Time { get; private set; }
public int Temperature { get; private set; }
public SensorData(DateTime time, int temperature)
{
Time = time;
Temperature = temperature;
}
}
Delay と ObserveOn
Rxのストリームはデータが発生したすぐ後にデータを流していくもの。
しかし、設計の都合上の理由からデータ発生直後ではなく少し時間をとってから流したい時がある。そのような場合に使われるのがDelay。
Delayの注意点は別スレッドに処理が移るという点である。以下の例ではその様子を確認する。
static void Main()
{
var form = new Form
{
Width = 300,
Height = 150,
Text = "演習"
};
var button = new Button
{
Text = "開始",
Width = 150,
Height = 50,
Top = 40,
Left = 70
};
form.Controls.Add(button);
Observable.FromEvent<EventHandler, EventArgs>(
onNext => (s, e) => onNext(e),
handler => button.Click += handler,
handler => button.Click -= handler)
.Do(_ =>
{
button.Enabled = false;
button.Text = "処理中...";
Console.WriteLine($"クリック時スレッドID: {Thread.CurrentThread.ManagedThreadId}");
})
.Delay(TimeSpan.FromSeconds(2))
.Subscribe(_ =>
{
Console.WriteLine($"Invoke前(Delay通過後)スレッドID: {Thread.CurrentThread.ManagedThreadId}");
button.Invoke((MethodInvoker)(() =>
{
button.Text = "完了";
button.Enabled = true;
Console.WriteLine($"完了時スレッドID: {Thread.CurrentThread.ManagedThreadId}");
}));
});
Application.Run(form);
}
今回は処理をInvokeで戻したが、Rxの書き方に即して書くならObserveOnが使える。
static void Main()
{
var form = new Form
{
Width = 300,
Height = 150,
Text = "演習"
};
var button = new Button
{
Text = "開始",
Width = 150,
Height = 50,
Top = 40,
Left = 70
};
form.Controls.Add(button);
Observable.FromEvent<EventHandler, EventArgs>(
onNext => (s, e) => onNext(e),
handler => button.Click += handler,
handler => button.Click -= handler)
.Do(_ =>
{
button.Enabled = false;
button.Text = "処理中...";
Console.WriteLine($"クリック時スレッドID: {Thread.CurrentThread.ManagedThreadId}");
})
.Delay(TimeSpan.FromSeconds(2))
.Do(_ => Console.WriteLine($"ObserveOn前(Delay通過後)スレッドID: {Thread.CurrentThread.ManagedThreadId}"))
.ObserveOn(SynchronizationContext.Current)
.Subscribe(_ =>
{
button.Text = "完了";
button.Enabled = true;
Console.WriteLine($"完了時スレッドID: {Thread.CurrentThread.ManagedThreadId}");
});
Application.Run(form);
}
ObserveOnはその後の処理を指定したスレッドで実行させる。
今後はSelectManyやThrottleについても勉強して、追記していけたら。










