【问题标题】:How to Merge two Observables so the result completes when the any of the Observables completes?如何合并两个 Observables,以便在任何 Observables 完成时结果完成?
【发布时间】:2011-02-03 15:17:53
【问题描述】:

我有这个代码:

var s1 = new Subject<Unit>();
var s2 = new Subject<Unit>();
var ss = s1.Merge(s2).Finally(() => Console.WriteLine("Finished!"));

ss.Subscribe(_ => Console.WriteLine("Next"));

s1.OnNext(new Unit());
s2.OnNext(new Unit());
s1.OnCompleted(); // I wish ss finished here.
s2.OnCompleted(); // Yet it does so here. =(

我已经使用 OnError(new OperationCanceledException()) 解决了我的问题,但我想要一个更好的解决方案(必须有一个组合器对吗?)。

【问题讨论】:

  • 很遗憾没有内置操作符,所以我写了一些扩展方法来解决这个问题。

标签: c# system.reactive combinators


【解决方案1】:

或者这个,也挺简洁的:

public static class Ext
{
    public static IObservable<T> MergeWithCompleteOnEither<T>(this IObservable<T> source, IObservable<T> right)
    {
        return Observable.CreateWithDisposable<T>(obs =>
        {
            var compositeDisposable = new CompositeDisposable();
            var subject = new Subject<T>();

            compositeDisposable.Add(subject.Subscribe(obs));
            compositeDisposable.Add(source.Subscribe(subject));
            compositeDisposable.Add(right.Subscribe(subject));


            return compositeDisposable;

        });     
    }
}

这使用了一个主题,该主题将确保在 CreateWithDisposable() 中仅将一个 OnCompleted 推送给观察者;

【讨论】:

  • 我建议使用AsyncSubject,否则您可能会遇到竞争条件。
  • 我比 Materialize() 更喜欢这个。 @Richard:你能解释一下它是如何引发比赛的吗?
  • 已更新。比赛是,当我订阅源和权限的主题时,OnComplete 和 subject.Subscribe(obs) 都会错过它。
  • 注意,我已经将 AsyncSubject 改回了 Subject
  • 注意:在最新的 Rx 版本中,Observable.CreateWithDisposable 已被重构为 Observable.Create 重载
【解决方案2】:

我建议将 onCompleted 事件转换为 onNext 事件并使用var ss = s1.Merge(s2).TakeUntil(s1ors2complete),而不是重写 Merge 以在任一流完成时完成,其中 s1ors2complete 在 s1 或 s2 结束时产生一个值。您也可以只链接 .TakeUntil(s1completes).TakeUntil(s2completes) 而不是创建 s1ors2complete。这种方法提供了比 MergeWithCompleteOnEither 扩展更好的组合,因为它可用于将任何“两个都完成时完成”运算符修改为“任何完成时完成”运算符。

关于如何将 onNext 事件转换为 onCompleted 事件,有几种方法可以做到这一点。 CompositeDisposable 方法听起来是个不错的方法,稍作搜索发现这个关于converting between onNext, onError, and onCompleted notifications 的有趣线程。我可能会使用xs.SkipWhile(_ =&gt; true).concat(Observable.Return(True)) 创建一个名为 ReturnTrueOnCompleted 的扩展方法,然后您的合并变为:

var s1ors2complete = s1.ReturnTrueOnCompleted().Amb(s2.ReturnTrueOnCompleted());
var ss = s1.Merge(s2).TakeUntil(s1ors2complete).Finally(() => Console.WriteLine("Finished!"));

您还可以考虑使用像 Zip 这样的运算符,当其中一个输入流完成时,automatically completes。

【讨论】:

  • P.S.以下是使用标准运算符组合多个流的不同方法的一个很好的概述:leecampbell.blogspot.com/2010/06/…
  • 不错的方法,我同意最好有一个操作员可以让我完成任何事情,而不仅仅是一个合并...谢谢!
【解决方案3】:

假设您不需要任何一个流的输出,您可以使用Amb 结合来自Materialize 的一些魔法:

var s1 = new Subject<Unit>();
var s2 = new Subject<Unit>();

var ss = Observable.Amb(
        s1.Materialize().Where(x => x.Kind == NotificationKind.OnCompleted), 
        s2.Materialize().Where(x => x.Kind == NotificationKind.OnCompleted)
    )
    .Finally(() => Console.WriteLine("Finished!"));

ss.Subscribe(_ => Console.WriteLine("Next"));

s1.OnNext(new Unit());
s2.OnNext(new Unit());

s1.OnCompleted(); // ss will finish here and s2 will be unsubscribed from

如果您需要这些值,您可以在两个主题上使用Do。

【讨论】:

  • 这只是带来了物化调用的成本,而且正如你所说,除了通过副作用 Do 之外,你无法获得这些值
【解决方案4】:

试试这个:

public static class Ext
{
    public static IObservable<T> MergeWithCompleteOnEither<T>(this IObservable<T> source, IObservable<T> right)
    {
        var completed = Observable.Throw<T>(new StreamCompletedException());

        return 
            source.Concat(completed)
            .Merge(right.Concat(completed))
            .Catch((StreamCompletedException ex) => Observable.Empty<T>());

    }

    private sealed class StreamCompletedException : Exception
    {
    }
}

它的作用是连接一个 IObservable,当源或正确的源完成时,该 IObservable 将引发异常。然后我们可以使用 Catch 扩展方法返回一个空的 Observable 以在任一完成时自动完成流。

【讨论】:

  • 这是我的第一次尝试,但我认为它不是很优雅(我使用了 OperationCanceledException)。
猜你喜欢
  • 2017-11-13
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-05-22
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多