- html - 出于某种原因,IE8 对我的 Sass 文件中继承的 html5 CSS 不友好?
- JMeter 在响应断言中使用 span 标签的问题
- html - 在 :hover and :active? 上具有不同效果的 CSS 动画
- html - 相对于居中的 html 内容固定的 CSS 重复背景?
我想合并 2 个可观察对象并保持顺序(可能基于选择器)。我还想对 observable 的来源施加反压。
因此选择器会选择其中一个项目通过可观察对象继续推进,而另一个项目将等待另一个项目也来进行比较。
Src1、Src2 和 Result 都是 IObservable<T>
类型.
Src1: { 1,3,6,8,9,10 }
Src2: { 2,4,5,7,11,12 }
Result: 1,2,3,4,5,6,7,8,9,10,11,12
Timeline:
Src1: -1---3----6------8----9-10
Src2: --2-----4---5-7----11---------12
Result: --1--2--3-4-5-6--7-8--9-10-11-12
这是否可以通过现有的 .net Rx 方法实现?
编辑:注意 2 个源可观察量保证是有序的。
示例测试:
var source1 = new List<int>() { 1, 4, 6, 7, 8, 10, 14 }.AsEnumerable();
var source2 = new List<int>() { 2, 3, 5, 9, 11, 12, 13, 15 }.AsEnumerable();
var src1 = source1.ToObservable();
var src2 = source2.ToObservable();
var res = src1.SortedMerge(src2, (a, b) =>
{
if (a <= b)
return a;
else
return b;
});
res.Subscribe((x) => Console.Write($"{x}, "));
期望结果:1,2,3,4,5,6,7,8,9,10,11,12,13,14,15
最佳答案
这很有趣。不得不稍微调整一下算法。它可以进一步改进。
假设:
streamA
, streamB
普通型T
.streamA[i] < streamA[i+1]
和 streamB[i] < stream[i+1]
.streamA[i]
之间有任何关系和 streamB[i]
.NotImplementedException
. 这种情况很容易处理,但我想避免歧义。min
对于类型 T
. 这是我使用的算法:
qA
和 qB
.streamA
获得商品时, 将其排入 qA
.streamB
获得商品时, 将其排入 qB
. qA
和 qB
,比较两个队列的顶部项目。移除并发出这两者的最小值。如果两个队列仍然非空,则重复。streamA
或 streamB
完成,转储队列的内容并终止。 注意:这确实是懒惰的,应该更改为转储,然后继续返回未完成的可观察对象。代码如下:
public static IObservable<T> SortedMerge<T>(this IObservable<T> source, IObservable<T> other)
{
return SortedMerge(source, other, (a, b) => Enumerable.Min(new[] { a, b}));
}
public static IObservable<T> SortedMerge<T>(this IObservable<T> source, IObservable<T> other, Func<T, T, T> min)
{
return source
.Select(i => (key: 1, value: i)).Materialize()
.Merge(other.Select(i => (key: 2, value: i)).Materialize())
.Scan((qA: ImmutableQueue<T>.Empty, qB: ImmutableQueue<T>.Empty, exception: (Exception)null, outputMessages: new List<T>()),
(state, message) =>
{
if (message.Kind == NotificationKind.OnNext)
{
var key = message.Value.key;
var value = message.Value.value;
var qA = state.qA;
var qB = state.qB;
if (key == 1)
qA = qA.Enqueue(value);
else
qB = qB.Enqueue(value);
var output = new List<T>();
while(!qA.IsEmpty && !qB.IsEmpty)
{
var aVal = qA.Peek();
var bVal = qB.Peek();
var minVal = min(aVal, bVal);
if(aVal.Equals(minVal) && bVal.Equals(minVal))
throw new NotImplementedException();
if(aVal.Equals(minVal))
{
output.Add(aVal);
qA = qA.Dequeue();
}
else
{
output.Add(bVal);
qB = qB.Dequeue();
}
}
return (qA, qB, null, output);
}
else if (message.Kind == NotificationKind.OnError)
{
return (state.qA, state.qB, message.Exception, new List<T>());
}
else //message.Kind == NotificationKind.OnCompleted
{
var output = state.qA.Concat(state.qB).ToList();
return (ImmutableQueue<T>.Empty, ImmutableQueue<T>.Empty, null, output);
}
})
.Publish(tuples => Observable.Merge(
tuples
.Where(t => t.outputMessages.Any() && (!t.qA.IsEmpty || !t.qB.IsEmpty))
.SelectMany(t => t.outputMessages
.Select(v => Notification.CreateOnNext<T>(v))
.ToObservable()
),
tuples
.Where(t => t.outputMessages.Any() && t.qA.IsEmpty && t.qB.IsEmpty)
.SelectMany(t => t.outputMessages
.Select(v => Notification.CreateOnNext<T>(v))
.ToObservable()
.Concat(Observable.Return(Notification.CreateOnCompleted<T>()))
),
tuples
.Where(t => t.exception != null)
.Select(t => Notification.CreateOnError<T>(t.exception))
))
.Dematerialize();
ImmutableQueue
来自System.Collections.Immutable
. Scan
需要跟踪状态。因为 OnCompleted
需要实现处理。诚然,这是一个复杂的解决方案,但我不确定是否有更清晰的以 Rx 为中心的方式。
如果您需要进一步说明,请告诉我。
关于c# - 合并两个 Observable 时保留排序,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/57238239/
一段时间后,我阅读了有关 RxJava concat 的内容,并决定测试一下我的理解力。但是我遇到了一些我不太理解的行为。 问题是,当我连接两个可观察对象时,根据我将它们传递给 Observable.
我正在使用来自数据库服务的数据实现自动完成: @Injectable() export class SchoolService { constructor(private db: AngularF
我正在尝试使用 RxJS 创建一个可观察的对象,它可以执行如图所示的操作。 获取一个值并等待一段固定的时间才能获得 下一个。 下一个将是该周期内发出的最后一个值 等等,跳过其余部分。 如果等待时间间隔
我有一个可观察对象和另一个提供的可观察对象改变 key 。我想构建一个在之间切换的可观察对象基于该键的对象中的可观察值。 示例: // Choose randomly between "up" or
我使用 protobuffers 在我的前端和我的 Dart 服务器之间进行通信。 那些对象没有实现 Observable . 我的 Dart 聚合物对象看起来像: @CustomTag('user-
在 java swing 项目中,我有一个模型类,它保存某个 JPanel 的状态。我需要使这些数据可供 View 使用。我认为有两种选择。有一个扩展 Observable 的类并将模型作为实例变量。
我想找到一种方法来检测观察者是否已完成使用我使用 Rx.Observable.create 创建的自定义可观察对象,以便自定义可观察对象可以结束它并正确地进行一些清理。 因此,我创建了一些测试代码,如
我正在尝试查询数据库。迭代结果列表,并为每一项再执行一个请求。在 rxjs 构建结束时,我有 Observable[]> 。但我需要Observable 。如何做到这一点? this.caseServ
我希望我的 api 上有一个方法返回 Observable> 但我希望该方法中的代码知道所有包含的 Observables 是否已完成,以便它可以关闭某些内容。最好的方法是什么? 更明确地说,我希望完
我有两个方法返回 Observable> firstAPI.getFirstInfo("1", "2"); secondApi.getSecondInfo(resultFirstObservable,
我有一个 Observable返回单个 Cursor实例(Observable)。我正在尝试利用 ContentObservable.fromCursor获取 onNext 中每个游标的行回调。 我想
我有两种返回 Observable 的方法: Observable firstObservable(); Observable secondObservable(String value); 对于第一
我正在尝试创建一个将用户数据作为 Observable 的函数,并使用来自第一个 observable 的数据从查询中添加/合并数据,然后将所有这些数据作为一个 observable 返回,我可以这样
我有一个 spec-compliant ECMAScript Observable ,具体来自 wonka library .我正在尝试将这种类型的 observable 转换为 rxjs 6 obs
为了简化问题,我在这里使用了数字和字符串。代码: const numbers$:Observable = of([1,2,3]); const strings: string[] = ["a","b"
对于我的 Android 应用程序,我需要一个 Observable 来聚合来自 7 个不同搜索的结果并作为一个集合发出。 对于最终发射,我选择了 ListMultimap其中 Content是搜索结
我正在使用改造 2.0.0-beta2 并且调试构建工作正常,但我在使用 Proguard 发布构建时遇到以下错误。 这是更新后的 logcat 错误。 11-17 18:23:22.751 1627
observer.throw(err) 和 observer.error(err) 有什么区别? 我正在使用 RxJS 版本“5.0.0-beta.12” var innerObservable =
我们有一种情况,对服务的方法调用返回一个 IObservable但我们的客户期望 IObservable .将 T1 转换为 T2 很简单。 Rx 中有什么允许这样做的吗? (即链接观察者) 我知道我
我陷入了如何将以下可观察类型转换/转换为我的目标类型的困境: 我有可观察的类型: Observable>> 我想将其转换为: Observable> 所以当我订阅它时,它会发出 List不是Obser
我是一名优秀的程序员,十分优秀!