極光日記

【Go】samber/roでリアクティブプログラミングする

Created:

samber/roでリアクティブプログラミング

はじめに

samber/roというライブラリを使ってGo言語でリアクティブプログラミングをしてみました。
使いこなせていませんが、使ってみて、気づいたことなどをメモします。

公式ドキュメントがわかりやすいため、網羅的な解説が必要な場合はそちらのほうが良いです。

基本

Observable

Observable型を作成するところから始まります。

ro.Just(10, 20, 30, 40) // ro.Observable[int]型

Operators

Observable型を元に何らかの処理(変換、フィルタリングなど)をすることになるでしょう。
Pipe関数を使うことで処理をつなげることができます。

pipeline := ro.Pipe1(
    ro.Just(10, 20, 30, 40),
    ro.Map(func(val int) int { return val * 2 }),
) 

このpipelineもro.Observable[int]型です。今回は変換処理がひとつだけだったのでPipe1を使いましたが、Pipe2,3,4〜も存在します。

なお、本稿ではPipeN関数ではなく、Pipe[First, Last any]を使った例文もあります。

pipeline := ro.Pipe[int, int](
    ro.Just(10, 20, 30, 40),
    ro.Map(func(val int) int { return val * 2 }),
)

これは上記の例文と同じ意味になります。こちらのほうがオペレーターの数が変わってもro.Pipe[int, int]を変えなくて良いメリットがあります。このため、検証などで試行錯誤する場合は便利なこともあります。
しかし、PipeN関数で書いたほうが型推論がされるため、[int, int]を省略でき、望ましいです。公式ドキュメントでもこちらの記法が推奨されています。

Subscription

リアクティブプログラミングの特徴ですが、前述のようにパイプラインを定義するだけでは何も起きません。
サブスクライバーを繋げる必要があります。

pipeline.Subscribe(ro.OnNext(func(val int) {
    fmt.Println(val) // 20, 40, 60, 80
}))
pipeline.Subscribe(ro.NewObserver(
    func(val int) {
        // onNext
        fmt.Println(val)
    },
    func(err error) {
        // onError
        fmt.Printf("Stream error: %v\n", err)
    },
    func() {
        // onCompleted
        fmt.Println("All streams completed!")
    },
))

もしくはブロッキング処理になってしまいますが、Collect関数もデバッグ時などに使えます。

result, err := ro.Collect(pipeline) // resultは[]int型。

エラー処理

まず、パイプラインの処理中でエラーが発生する場合、MapではなくMapErrを使う必要があります。
公式ドキュメント記載の例文を少し変更したものですが、MapErrが使われ、これにはfunc(i int) (int, error)のようにエラーを返す関数が渡されています。

obs := ro.Pipe[int, int](
    ro.Just(1, 2, 3),
    ro.MapErr(func(i int) (int, error) {
        if i == 2 {
            return 0, errors.New("number 2 is not allowed")
        }
        return i * 2, nil
    }),
    ro.Catch(func(err error) ro.Observable[int] {
        fmt.Printf("Error: %v\n", err)
        return ro.Just(99) // Fallback value
    }),
)

result, err := ro.Collect(obs)
fmt.Println("Result:", result) // Result: [2 99]
fmt.Println("Error:", err) // Error: <nil>

さらに、この例文ではro.Catchが続いています。
その前にエラーが発生した時にro.Catchの中の処理が走ります。
ただし、このCatchでFallback値を返したとしても処理はそこで終了してしまうようです。コメントにある通り、結果はResult: [2 99]となり、3が処理されていません。
個人的にはエラーが発生しても処理を続けたい場合は、次のようにエラーを返さず、Map関数でFallback値を返すようにしました。

obs := ro.Pipe[int, int](
    ro.Just(1, 2, 3),
    ro.Map(func(i int) int {
        if i == 2 {
            // エラー時のフォールバック
            return 99
        }
        return i * 2
    }),
)

Empty

途中の処理でro.Emptyを返すこともできます。
Emptyを返された場合、それ以降の処理はスキップされます。下記の例文ではコメントに記載されている通り、2はスキップされ、1, 3は処理されています。
前述のエラー処理に関連して、エラー発生時はEmptyを返すのもありかもしれません。
なお、この例のようにプレーンな型ではなくro.Observable型を返したい場合はMapではなく、FlatMapを使います。

obs := ro.Pipe[int, int](
    ro.Just(1, 2, 3),
    ro.FlatMap(func(i int) ro.Observable[int] {
        if i == 2 {
            fmt.Println("number 2 will be empty")
            return ro.Empty[int]()
        }
        return ro.Just(i * 2)
    }),
    ro.Map(func(i int) int {
        return i * 10
    }),
)

result, err := ro.Collect(obs)
fmt.Println("Result:", result) // Result: [20 60]
fmt.Println("Error:", err) // Error: <nil>

2つの値を返したい

リアクティブプログラミングでMapを繋げていく場合、シンプルに入力を変換するだけなら良いのですが、入力値+αを後続に送りたくなることがあります。
その場合、新しく構造体などを作成すれば良いのですが、値の受け渡しのためだけに構造体を作成するのは面倒なこともあります。
その場合、roと同じ作者が作成しているライブラリであるloのTupleが使えます(loの導入方法など)。以下が例文です。

pipeline := ro.Pipe2(
    ro.Just("apple", "banana", "cherry", "date"),
    ro.Map(func(word string) lo.Tuple2[string, int] {
        l := len(word) // 実際はDBクエリをするなど
        return lo.T2(word, l)
    }),
    ro.Map(func(pair lo.Tuple2[string, int]) int {
        fmt.Printf("word: %s, len: %d\n", pair.A, pair.B)
        return pair.B
    }),
)

ただし、多用するとA、Bがそれぞれ何なのかがわからなくなります。 シンプルな場合にとどめ、それ以外は構造体を用意したほうが良いでしょう。

テスト

次のような、ro.Observable型を返す関数をテストしたい場合、

func example() ro.Observable[int] {
	pipeline := ro.Pipe[int, int](
		ro.Just(10, 20, 30, 40),
		ro.Map(func(val int) int { return val * 2 }),
	)

	return pipeline
}

ro.Collectを使えば、結果を取り出せるので、それを使えば、普通にテストを書くこともできます。

func TestExample(t *testing.T) {
	got, _ := ro.Collect(example())

	want := []int{20, 40, 60, 80}
	if !slices.Equal(got, want) {
		t.Errorf("example() = %v, want %v", got, want)
	}
}

ただし、roにはテスト用の機能も用意されており、次のようにも書けます。

import rotesting "github.com/samber/ro/testing"

func TestExample(t *testing.T) {
	rotesting.Assert[int](t).
		Source(example()).
		ExpectNext(20).
		ExpectNext(40).
		ExpectNext(60).
		ExpectNext(80).
		// ExpectNextSeq(20, 40, 60, 80). // ExpectNextSeqを使ってまとめて書くこともできる
		ExpectComplete().
		Verify()
}

解説:

  • ExpectComplete
    • ストリームが正常終了したことを確認する。これがないと予想外に100などが流れてきても失敗にならない。
  • Verify
    • ストリームをサブスクライブする。これがないとストリームが流れ始めず、何も起きない。

この方がリアクティブプログラミングらしいテストになるかもしれません。
また、このようなシンプルなテストをする以外にも色々な機能があるようです。

参考: https://ro.samber.dev/docs/testing#Simple-Assertion

どこで使うべきか?

本稿に含まれる例文のような処理をするならば普通にfor文を使った方がわかりやすく、リアクティブプログラミングするメリットはありません。
公式ドキュメントには

  • Real-time data processing (WebSocket events, sensor data)
  • User interface events (clicks, keystrokes, form inputs)
  • API response handling (with retry, timeout, and caching)
  • Data processing with transformation, aggregation and enrichment
  • Event-driven patterns

とあります。正直あまり具体的な使いどころは思いつかないのですが、リッチなCLIアプリケーションは使いどころの一つかもしれません。
例えば、claude codeやcodexのようなTUIならば、文字を入力、Enterを押下して送信、ショートカットキー操作も可能といった機能を実装する必要があります。これはまさに「User interface events (clicks, keystrokes, form inputs)」となります。
一方、Goが使われる事の多い、典型的なWebアプリケーションの場合、JSONを受信→DBをクエリ→DBに登録→結果を返信といった逐次処理をするだけならばリアクティブプログラミングをする必要は少ないと思います。roはリトライなどの処理も書けますが、もっと普通に実装した方が多くの場合シンプルになるでしょう。

また、Goの文法がリアクティブプログラミングに向いていないとも思います。
リアクティブプログラミングならば

ro.Just(10, 20, 30, 40).Map(i => i*2).FlatMap(x => foo(x))

などのように書きたいものですが、Goではこのようにできません。ラムダ式がなく、errorを含めて2値を返すことがあるため、メソッドチェーンも難しいです。そのため、以下のようにPipe関数を使って書きますが、後からOperatorを増やしたくなった時に不便です。

obs := ro.Pipe3(
    source,
    ro.Filter(predicate),
    ro.Map(transformer),
    ro.Take[int](10),
)