MS_ReactiveExtensions - NetDevInfraWGinOSSConsortium/NetDevInfraWiki GitHub Wiki

Reactive ExtensionsRx

抂芁

  • Observer パタヌンを実装するフレヌムワヌクだが、

  • 特に、LINQ to Events もしくは LINQ to Asynchronous ず衚珟できる。

  • これは、芁するに、

    • むベントや非同期凊理を LINQ っぜく扱える。
    • 非同期 / むベント / 時間に関する凊理を LINQ 的に簡朔か぀宣蚀的に蚘述できる。

ず蚀うトコロを意味する。

移行メモ衚蚘: 原文の「LINQ to Asynchronus」は
Asynchronous の綎りの誀りず刀断し修正した。

補足IEnumerable ず IObservable は双察: Rx を理解する鍵は、
IObservable<T> が IEnumerable<T> の裏返しである、ずいう点にある。

IEnumerable<T>Pull IObservable<T>Push
䞻導暩 受け手が MoveNext() で取りに行く 送り手が OnNext() で抌し蟌む
終了 MoveNext() が false OnCompleted()
ç•°åžž 䟋倖が throw される OnError()
挔算子 Where / Select / 
 同じ名前の挔算子が䜿える

「倀の䞊び」であるこずは同じなので、
LINQ の挔算子がそのたた通甚する。
これが「むベントを LINQ で扱える」ずいう発想の正䜓である。

// 「クリックが 300ms 以内に 2 回」 ダブルクリック、を宣蚀的に曞く
clicks.Buffer(clicks.Throttle(TimeSpan.FromMilliseconds(300)))
      .Where(x => x.Count == 2)
      .Subscribe(_ => Console.WriteLine("double click"));

時間を第䞀玚に扱える点が、Rx が今も䟡倀を持぀理由である。

経緯

  • Silverlight Toolkit に "System.Reactive.dll" が同梱される。

  • Reactive Framework ---> Reactive Extensions ず名称倉曎。

  • DevLabs昔の MS コミュニティサむトでプロゞェクト公開
    ゜ヌスコヌドは恐らくCodePlex で公開されおいた。

  • 䜕床も API の消滅・远加を繰り返す。

  • JavaScript 版が登堎RxJS

  • .NET Framework 4 SP1 に暙準搭茉

  • さたざたな開発蚀語に移怍され、幅広く䜿われおいる。

    • JavaScript 甚の RxJS
    • Java / Android 甚の RxJava
    • Swift 甚の RxSwift
    • Unity 甚の UniRx

補足その埌最新化: 原文以降の動きを補う。

時期 内容
2016 頃 ReactiveX ずしお蚀語暪断の仕様・サむトに集玄
Rx.NET が .NET Foundation 配䞋でメンテナンスされる
2019 IAsyncEnumerable<T> が C# 8.0 / .NET Core 3.0 に暙準搭茉
珟圚 System.Reactive パッケヌゞずしお提䟛Rx.NET リポゞトリ
Unity UniRx → R3 / UniTask に䞖代亀代が進む

重芁な倉化は IAsyncEnumerable<T> の登堎である。
「非同期に届く倀の䞊びを扱う」ずいう Rx の甚途の䞀郚が、
蚀語機胜await foreachで玠盎に曞けるようになった。

// 昔は Rx が必芁だった凊理が、暙準構文で曞ける
await foreach (var item in GetItemsAsync())
    Console.WriteLine(item);

珟圚の䜿い分けの目安は次の通り。

やりたいこず 適した手段
1 回の非同期凊理 async/await
非同期に届く有限の䞊びPull IAsyncEnumerable<T>
時間・むベントの合成Push Rx
UI のむベント凊理・デバりンス Rx

Rx は「䞇胜の非同期基盀」から
**「むベントず時間の合成に特化した道具」**ぞず䜍眮づけが定たった、
ず理解するのが実態に近い。

ナヌスケヌス

Observer パタヌンの実装

LINQフィルタ、倉換、集蚈、合成

通知に察しお、LINQフィルタ、倉換、集蚈、合成を適甚できる。

むベントの LINQ 化

むベントに Observer パタヌンを適甚しお、
LINQフィルタ、倉換、集蚈、合成を適甚できる。

非同期の LINQ 化

  • むベントに Observer パタヌンを適甚しお、
    LINQフィルタ、倉換、集蚈、合成を適甚できる。

  • 具䜓的には、Observer を合成しおネストを凊理が可胜。

詳现

Observer パタヌンで構成される。

  • 䜕らかのProviderIObservable<T> などで監芖を行い、
  • Providerから状態の Push通知、発行、むベントを、
  • SubscriberIObserver<T> などが受け取る。

Subscriber

IObserver<T> を実装したクラスを定矩する。

Provider

  • IObservable<T> を実装したクラスを定矩する。
  • クラスはラムダ匏を䜿甚したメ゜ッドチェヌンで実装できる。

開始

  • Create などの生成メ゜ッドで生成、
# メ゜ッド 抂芁
1 Create 監芖内容を定矩する。
2 Return 枡した倀を単玔に通知する。
3 Range 指定した範囲の倀を通知する。
4 Repeat 第 1 匕数で枡したデヌタを、
第 2 匕数で指定した回数繰り返しお通知する。
5 Generate for 文的に以䞋を定矩する。
・初期倀
・継続刀定デリゲヌト
・むンクリメント・デリゲヌト
・通知する倀を生成するデリゲヌト
6 Case 耇数甚意された IObservable の䞭から、どれか぀を遞択。
・IObservable 配列
・IObservable 遞択デリゲヌト
7 Throw 䟋倖を通知する。
8 Start 枡した倀を非同期で通知する。
9 Defer Provider のファクトリを実装。

移行メモ誀字: Range の説明「正指定した範囲の倀」は
指定した範囲の倀の誀字ず刀断し修正した。

  • Create で生成する堎合、以䞋のメ゜ッドを䜿甚しお監芖内容を定矩する。
# メ゜ッド 抂芁
1 OnNext 新しい倀が発生したこずを通知する。
2 OnError ゚ラヌが発生し、異垞終了したこずを通知する。
3 OnCompleted 正垞に終了したこずを通知する。

補足Rx の文法芏則: この 3 ぀には厳密な順序の玄束がある。

OnNext* (OnError | OnCompleted)?

぀たり、

  • OnNext は 0 回以䞊流れる
  • OnError たたは OnCompleted は どちらか䞀方が、最埌に 1 回だけ
  • 終了通知の埌には、䜕も流れおはならない

自分で IObservable<T> を実装する堎合、この芏則を守る責任は
実装偎にある。Observable.Create を䜿えば倧郚分は担保される。

  • 必芁に応じお、LINQ メ゜ッドを䜿甚し通知内容を操䜜できる。

  • Subscribe で、䞊蚘の其々のメ゜ッドからの
    Push通知、発行、むベントを受け取るSubscriberを蚭定する。

終了

  • Provider が OnCompleted を呌び出し、正垞終了。
  • Provider が OnError を呌び出し、異垞終了。
  • Subscriber偎で Dispose する。

補足賌読解陀の忘れがリヌクになる: Subscribe は
IDisposable を返す。これを捚おるず賌読が残り続け、
Provider が Subscriber を参照し続けるためメモリ リヌクになる。

// ① using / Dispose で明瀺的に解陀
using var sub = observable.Subscribe(...);

// ② CompositeDisposable にたずめお、画面砎棄時に䞀括解陀
_disposables.Add(observable.Subscribe(...));

// ③ 寿呜を別のストリヌムに委ねる
observable.TakeUntil(_closed).Subscribe(...);

OnError で賌読が終了する点も重芁で、
゚ラヌ埌も流し続けたい堎合は Retry / Catch を挟む必芁がある。

分類

倧きく分けお Hot / Cold に分類される。

  • Cold

    • Subscriberずはの関係
    • Subscriberが居なければ動䜜しない。
    • 監芖が終わった埌は OnCompleted で自分で終了する。
    • Observable.Interval など䟋倖もある。
  • Hot

    • Subscriberずはの関係
    • Subscriberがいなくおも動䜜し続ける。
    • 監芖が終わったかどうかSubscriberが刀断しお Dispose。

補足Hot / Cold は Rx 最倧のハマりどころ: 原文が
埌述のポむントで「抂念を掎む」ず匷調しおいる通り、
ここが最も事故が起きる箇所である。

Cold Hot
倀の生成 賌読のたびに最初から生成 1 本の流れを共有
䟋 Observable.Range、HTTP リク゚スト むベント、Subject、株䟡
賌読前の倀 倱われない賌読時に始たる 賌読前の倀は受け取れない

兞型的な事故:

var src = Observable.Create<int>(o => { Console.WriteLine("実行"); ... });
src.Subscribe(...);   // 「実行」
src.Subscribe(...);   // 「実行」← 2 回動いおしたう

HTTP リク゚ストなど副䜜甚のある Cold を耇数箇所で賌読するず、
リク゚ストが賌読数だけ飛ぶ。

察凊は Cold を Hot に倉換するこず。

var shared = src.Publish().RefCount();  // 賌読者間で 1 本を共有

Subject (Subscriber/Provider)

  • IObserver<T> ず IObservable<T> を実装したクラスを定矩できる。
  • Subject を䜿甚すれば、Subscriber をラムダ匏で定矩できる。

補足Subject の皮類ず、䜿いすぎぞの泚意:

皮類 挙動
Subject<T> 賌読埌に流れた倀のみ受け取る
BehaviorSubject<T> 最新の 1 件を賌読時に即座に受け取る状態の衚珟に適する
ReplaySubject<T> 過去 n 件を賌読時に再生
AsyncSubject<T> 完了時に最埌の 1 件だけ流すTask に近い

ただし、Subject を安易に䜿うず Rx の宣蚀的な利点が倱われる
手続き的に OnNext を呌ぶコヌドが散らばる。
可胜なら Observable.FromEvent や挔算子で組み立お、
Subject は倖郚ずの境界に限定するのが定石である。

その他

合成

Provider の合成

  • Merge
  • SelectMany
  • Switch

倀の合成

  • Concat
  • Zip
  • Amb
  • CombineLatest

補足実務で䜿う頻床が高いもの: 数癟ある挔算子のうち、
実際によく䜿うのは限られる。

挔算子 甹途
ThrottleRxJS の debounce 入力が止たっおから実行むンクリメンタル怜玢
Sample 䞀定間隔で最新倀を間匕く
Switch 新しい芁求が来たら、前の凊理を捚おる怜玢の競合防止
CombineLatest 耇数の最新倀を組み合わせるフォヌムの劥圓性刀定
Retry / Catch ゚ラヌ時の再詊行
DistinctUntilChanged 同じ倀の連続を無芖

Throttle + Switch の組み合わせは、
「入力のたびに怜玢 API を叩き、叀い結果が埌から返っお䞊曞きする」
ずいう叀兞的なバグを、宣蚀的に解決する。

textChanged
  .Throttle(TimeSpan.FromMilliseconds(300))   // 入力が萜ち着いおから
  .DistinctUntilChanged()                     // 同じ語なら投げない
  .Select(q => SearchAsync(q).ToObservable())
  .Switch()                                   // 叀い怜玢は砎棄
  .Subscribe(UpdateUI);

Merge ず Concat の違いにも泚意。
Merge は䞊行、Concat は前が終わっおから次である。

倉換

むベント

  • FromEvent メ゜ッドで、
    むベント・クラスを IObservable に倉換できる。

  • LINQ メ゜ッドを䜿甚し通知内容を操䜜できる。

    • 具䜓的には、ドラッグ操䜜を簡単に蚘述できる。
  • 参考

補足ドラッグが Rx の代衚䟋である理由: ドラッグは
「MouseDown しおから MouseUp するたでの MouseMove」であり、
むベントの時間的な組み合わせそのものである。

var drag = from down in mouseDown
           from move in mouseMove.TakeUntil(mouseUp)
           select move.Location;

手続き的に曞くず「ドラッグ䞭フラグ」ず開始座暙の状態管理が芁るが、
Rx では状態倉数を持たずに衚珟できる。

非同期

  • FromAsyncPattern メ゜ッドで、
    BeginInvoke/EndInvoke を IObservable に倉換できる。

  • LINQ メ゜ッドを䜿甚し通知内容を操䜜できる。

    • 具䜓的には、
      • SelectMany を組み合わせ、
      • 非同期のネストを簡単に蚘述できる。
  • 参考

補足珟圚は async/await が担う領域最新化: この節が扱う
APMBegin/Endのネスト解消は、
Rx が登堎した圓時.NET 4 以前における最倧の動機の䞀぀だった。

しかし C# 5.0 の async/await により、
「非同期のネスト」は蚀語レベルで解決された。

// か぀お Rx / SelectMany で解いおいた問題
var a = await GetAAsync();
var b = await GetBAsync(a);

珟圚の Rx ず Task の盞互倉換は次の通り。

倉換 メ゜ッド
Task → IObservable .ToObservable()
IObservable → Task await observable最埌の倀、.ToTask()
IObservable → IAsyncEnumerable .ToAsyncEnumerable()

なお FromAsyncPattern は
APM 自䜓が非掚奚Begin/End は珟圚ほが䜿われないのため、
新芏コヌドでは Observable.FromAsync を䜿う。

ポむント

Hot ず Cold の抂念を掎む

Scheduler が実行スレッドを決定する

補足Scheduler の芁点: Rx ではどのスレッドで動くかを
挔算子ではなく Scheduler が決める。

メ゜ッド 䜕を倉えるか
ObserveOn 通知を受け取る䞋流が動くスレッド
SubscribeOn **賌読凊理䞊流の開始**が動くスレッド

UI アプリでは、
**「重い凊理はバックグラりンド、UI 曎新は UI スレッド」**を
この 2 ぀で宣蚀的に指定できる。

source
  .SubscribeOn(TaskPoolScheduler.Default)      // 取埗はプヌルで
  .Select(Heavy)
  .ObserveOn(DispatcherScheduler.Current)      // 衚瀺は UI スレッドで
  .Subscribe(UpdateUI);

既定では、Rx はスレッドを勝手に切り替えない
OnNext を呌んだスレッドがそのたた䞋流を実行する点に泚意。

Subject でテストする

補足テストの本呜は TestScheduler: 原文の蚀う
「Subject でテストする」任意のタむミングで倀を流し蟌むに加え、
Rx には TestSchedulerMicrosoft.Reactive.Testingがある。

仮想時間を進められるため、
Throttle(5分) のような凊理を
実際に 5 分埅たずにテストできる。

var s = new TestScheduler();
var result = s.Start(() => source.Throttle(TimeSpan.FromMinutes(5), s));
// 5 分を「即座に」進めお怜蚌できる

時間に䟝存する凊理を確実にテストできるずいう点は、
手続き的な実装に察する Rx の明確な優䜍性である。

事件

2020/03/09 ちょっずした事件が起きたらしい。

FF

ファクト・ファむンディング

  • Qiita

    • UniRxを䜿った開発をしお思ったこず
      https://qiita.com/nodead/items/781ad247c127af47a1a2

    • プログラム開発者栌差の話をしよう
      https://qiita.com/Rwf-9DH3/items/fe65756c0485f095cc98

      「なるほど完璧な䜜戊っスね―――ッ
      䞍可胜だずいう点に目を぀ぶればよぉ」
      (匕甚ゞョゞョの奇劙な冒険 第4郚

感想

  • むむね。

  • SI っぜくなっお来た。

  • dis っおは無い。

「そう蚀う事も考える必芁が出おきた。」

...ず蚀う事で。

補足この「事件」の論点: 参照先は、
「Rx は匷力だが、チヌム党䜓が䜿いこなせるずは限らない」
ずいう、技術遞定における習熟床栌差の話題である。

Rx は孊習曲線が急で、

  • Hot / Cold を理解しおいないず、副䜜甚が倚重に走る
  • 賌読解陀を忘れるずリヌクする
  • 挔算子が数癟あり、レビュヌで劥圓性を刀断しづらい
  • スタック トレヌスが远いにくく、デバッグが難しい

ずいった性質を持぀。原文の「SI っぜくなっお来た」ずいう感想は、
個人の技量に䟝存する技術を、倚人数の開発に持ち蟌む難しさを
指したものず読める。

珟圚の劥圓な刀断ずしおは、

  • 単発の非同期は async/await、
  • 有限の非同期列は IAsyncEnumerable<T>、
  • 本圓に「時間ずむベントの合成」が芁る箇所に限っお Rx、

ず切り分けるのが、この論点ぞの実務的な回答になる。

参考

移行メモリポゞトリ移管: 原文の
Reactive-Extensions/Rx.NET は、珟圚 dotnet/reactive に移管
.NET Foundation 配䞋されおいるため、URL を曎新した。

IT

Build Insider

Qiita

neue cc

xin9le.net

かずきのBlog@hatena

present

Reactive Extensions 入門

Muhammad Rehan Saeed

Reactive Extensions (Rx)


Tags: 移行, .NET開発

⚠ **GitHub.com Fallback** ⚠