【发布时间】:2021-02-26 15:58:19
【问题描述】:
我搜索了一个副本,但没有找到。我拥有的是一个嵌套的 observable IObservable<IObservable<T>>,我想将它展平为 IObservable<T>。我不想使用 Concat 运算符,因为它会延迟对每个内部 observable 的订阅,直到前一个 observable 完成。这是一个问题,因为内部 observable 很冷,我希望它们在外部 observable 发出 T 值后立即开始发出。我也不想使用Merge 运算符,因为它会打乱发出值的顺序。下面的大理石图显示了Merge 运算符的有问题的(就我而言)行为,以及理想的合并行为。
Stream of observables: +----1------2-----3----|
Observable-1 : +--A-----------------B-------|
Observable-2 : +---C---------------------D------|
Observable-3 : +--E--------------------F-------|
Merge (undesirable) : +-------A-------C----E----B-----------D---F-------|
Desirable merging : +-------A-----------------B-------C---D------EF---|
Observable-1 发出的所有值都应该在 Observable-2 发出的任何值之前。 Observable-2 和 Observable-3 也应如此。
我喜欢Merge 运算符的原因是它允许配置对内部可观察对象的最大并发订阅。我想使用我正在尝试实现的自定义 MergeOrdered 运算符保留此功能。这是我正在构建的方法:
public static IObservable<T> MergeOrdered<T>(
this IObservable<IObservable<T>> source,
int maximumConcurrency = Int32.MaxValue)
{
return source.Merge(maximumConcurrency); // How to make it ordered?
}
这是一个用法示例:
var source = Observable
.Interval(TimeSpan.FromMilliseconds(300))
.Take(4)
.Select(x => Observable
.Interval(TimeSpan.FromMilliseconds(200))
.Select(y => $"{x + 1}-{(char)(65 + y)}")
.Take(3));
var results = await source.MergeOrdered(2).ToArray();
Console.WriteLine($"Results: {String.Join(", ", results)}");
输出(不良):
Results: 1-A, 1-B, 2-A, 1-C, 2-B, 3-A, 2-C, 3-B, 4-A, 3-C, 4-B, 4-C
理想的输出是:
Results: 1-A, 1-B, 1-C, 2-A, 2-B, 2-C, 3-A, 3-B, 3-C, 4-A, 4-B, 4-C
澄清:关于值的顺序,值本身是无关紧要的。重要的是它们起源的内部序列的顺序,以及它们在该序列中的位置。第一个内部序列中的所有值应首先发出(按其原始顺序),然后是第二个内部序列中的所有值,然后是第三个内部序列中的所有值,依此类推。
【问题讨论】:
-
This 是我能找到的最接近的“重复”。它有点相关,但有一些具体的细微差别,使它对我的问题既宽又窄。
-
有趣,你有这个可以玩,还有源代码,也许可以让你走上正轨:rxmarbles.com
-
您是否考虑过从源中的冷的 observable 创建热的 observable,然后只使用
Concat? -
@Kamushek 是的,我也有这个想法,通过使用
Publish运算符。但是我正在失去价值,而且我还没有设法将它与maximumConcurrency功能结合起来。 -
@aybe rxmarbles.com 是否允许创建请求的
MergeOrdered运算符的这个问题中显示的复杂弹珠图? AFAICS 我只能选择一个预定义的运算符,然后用鼠标左右移动图表项目符号。
标签: c# concurrency system.reactive rx.net