如何通过响应式扩展将一个事件拆分为多个事件?

rha*_*der 3 arrays split stream reactive-programming system.reactive

如何在响应式扩展流中处理单个事件并将其拆分为同一个流中的多个事件?

我有一个序列来检索json数据,这是一个顶层的数组.在json数据被解码的那一点上,我想获取该数组中的每个元素并继续沿着流传递这些元素.

这是一个我希望存在的虚构函数的例子(但名字较短!).它是在Python中,但我认为它很简单,它应该对其他Rx程序员清晰.

# Smallest possible example
from rx import Observable
import requests
stream = Observable.just('https://api.github.com/users')
stream.map(requests.get) \
      .map(lambda raw: raw.json()) \
      .SPLIT_ARRAY_INTO_SEPARATE_EVENTS() \
      .subscribe(print)
Run Code Online (Sandbox Code Playgroud)

换句话说,我想像这样进行转换:

From:
# --[a,b,c]--[d,e]--|->
To:
# --a-b-c-----d-e---|->
Run Code Online (Sandbox Code Playgroud)

Mar*_*age 5

您可以使用SelectMany运营商:

stream.SelectMany(arr => arr)
Run Code Online (Sandbox Code Playgroud)

这将"平坦化"您的事件流,就像C#LINQ SelectMany运算符可用于展平序列序列一样.