1
1

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

Reactive Extensions の備忘録

1
Last updated at Posted at 2026-10-04

はじめに

誤りについては優しく指摘していただけますと幸いです(まだまだ理解できていないので誤りだらけかと思います)。

Reactive Extensions (Rx) とは何か

一言でいうと、状態変化に応じた処理が書きやすいことが利点。
対象を観測し、その状態が変化すると観測側に通知が行くシステム。

いつ何回発生するか不明なイベントがあり、
それについて処理を行うシチュエーションを考える。

LINQのように処理する側 = データを要求する側に主導権があるとスレッドをブロックしてしまう。 これは同期的に値を返すことを前提としているために引き起こされる。(Pull型)
一方Rxではイベント側 = データを提供する側に主導権があるため、スレッドをブロックせずに済む。 これは、値が来たら提供側から届けるという契約になっているため。(Push型)

LINQ
        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;
        }

実行例
LINQ.gif

Rx
    static void Main()
    {
        Console.WriteLine("開始");

        // 3秒後に「1」というデータを通知し、表示する
        Observable.Timer(TimeSpan.FromSeconds(3))
            .Subscribe(x => Console.WriteLine("1"));

        // スレッドがブロックされないため、3秒待たずに「終了」と表示される。
        Console.WriteLine("終了");
        Console.ReadKey();
    }

実行例

Rx.gif

似たようなことがAsync/Awaitでも可能だが、あちらは連続して起こる処理には向かない。
「状態変化」はいつ起きるか分からない上に何度も発生しうる。
そんな時にはRxが威力を発揮する。

Rxは通知を便利に扱うライブラリ。
なのでいたずらに使うと、値が可変になってしまい クラスの安定性が低下する。
利用する側の実装難易度が上がってしまうというリスクも持っている。

使い分けまとめ

LINQ Async/Await Rx
DBから取得して
加工
WebAPIを1回呼び出す UIのイベント処理
メモリ上の
配列の集計
順序だった非同期処理 リアルタイム通知の処理
CSV/JSONの
読み込み
並行して処理し、
全て終了するのを待つ
WebSocketのプッシュ通知

用語のイメージ

用語 イメージ
ストリーム 川の流れ
Observable 川の設計・水源
Observer 川下で水を受け取る人
subscribe 水門の開放
subject 誰でも水を注ぎ・汲める池
実データ 水

Where/Select/Subscribe

基本的なRxのメソッドの具体例。

Where/Select/Subscribe
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();
    }
}

実行例
スクリーンショット 2026-10-03 223434.png

Where

LINQのWhereと同様に、流れて来るデータをフィルタリングする。

Select

LINQのSelectと同様に、流れて来るデータを加工する。

Subscribe

データを流し始めるという宣言。これが呼ばれた部分からデータが流れるので、これ以前に流されたデータは渡されない。戻り値はIDisposable。
IObservable<T>が呼び出すことができ、 引数にはIObserver<T>が入る。

  • onNextには変化が通知された場合に行う処理が書かれており、
    IObservableはIObserver(つまり書かれた処理)にデータを渡し、処理を行う。
  • onCompletedはデータが流れ終わったときに行われる処理が書かれており、
    これ以降に流れたデータは通知されない。
  • onErrorというものもあり、こちらはエラーの通知が来た時に行う処理が書かれる。

GroupBy

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}"));
            });
    }

実行例
スクリーンショット 2026-10-03 232303.png

GropuByは

  1. 流れてきたストリームの分割
  2. 分割されたグループに値を流す

という2つのことを行っていると考えている。

つまり、sourceで1~10までの数値というストリームが作られてGroupByに流れて来る。

  1. まず1が流れてくると奇数のグループとして IGroupedObservableを生成する
  2. その後、そのグループには「グループのKeyと流れてきた値を出力する」という処理を登録して1を渡す。すると、グループのKey(今回は2で割ったあまり) と渡された1が表示される
  3. 次に2が渡された際の処理も同様。3が渡された場合は奇数グループになるが、そのグループは既に存在するため新規の作成は行われない。その工程はスキップされ、そのままストリームが流れていく (データが渡される) 形になる

FromEvent

WinFormsなどの既存ライブラリは通知をeventで提供している。
しかし、WhereやSelectといった演算子はIObservable<T>にしか使えない。
既存のイベントをRxで扱う手段として、FromEventが挙げられる。

これにより、eventをIObservable<T>という「値」として扱えるようになり、
Whereや後述のMergeに渡せるようになる。
また、eventはKeyEventHandler、EventHandler、MouseEventHandlerなどイベントごとに異なるシグネチャを持っているが、
これを統一的に扱えるようになることもメリットとして挙げられる。

以下のサンプルコードではFormが表示され、
Enterキーを押した時にだけコンソールにReturnと表示される。

FromEvent
        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);

最初にこれを見せられても良く分からなかったので、
以下にラムダ式を減らして書いたバージョンも記載する。

FromEvent(ラムダ式を減らしたバージョン)
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;
    }
}

実行例
スクリーンショット 2026-10-04 100116.png

現時点での理解を記しておくが、読みづらいので後々読みやすい形にする予定。

  • 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。

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);
    }

実行例
スクリーンショット 2026-10-04 171158.png

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回ずつ値を出すまでは何もしないのが特徴。

CombineLatest
    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);
    }

実行例
スクリーンショット 2026-10-04 173747.png

DistinctUntilChanged と Distinct

直前と同じ値の場合は通知されない。
「変化」のみを伝えることで、無駄な通知を省くことができる。

DistinctUntilChanged
    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();
        }
    }

実行例
スクリーンショット 2026-10-10 092011.png

類似したものにDistinctがあるが、こちらは今までに通知していない値のみを通知する。
今までに1度でも出ている値は通知されない。
挙動が異なるので注意。

下のコード例では気温にのみ基づいて同じ値か否かを比較するため、
data => data.Temperatureと記載している。

Distinct
    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;
        }
    }

出力例
スクリーンショット 2026-10-10 094025.png

Delay と ObserveOn

Rxのストリームはデータが発生したすぐ後にデータを流していくもの。
しかし、設計の都合上の理由からデータ発生直後ではなく少し時間をとってから流したい時がある。そのような場合に使われるのがDelay。

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);
    }

実行例
Delay.gif

今回は処理をInvokeで戻したが、Rxの書き方に即して書くならObserveOnが使える。

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);
    }

実行例
Delay2.gif

ObserveOnはその後の処理を指定したスレッドで実行させる。

今後はSelectManyやThrottleについても勉強して、追記していけたら。

1
1
0

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
1
1

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?