【问题标题】:Reactive Extensions Synchronous Subscription反应式扩展同步订阅
【发布时间】:2014-07-18 04:28:00
【问题描述】:

有人可以帮我对 IObserver 进行同步订阅,这样调用方法就会阻塞,直到订阅完成。 例如:

出版商

public static class Publisher {
public static IObservable<string> NonBlocking()
    {
        return Observable.Create<string>(
            observable =>
            {
                Task.Run(() =>
                {
                    observable.OnNext("a");
                    Thread.Sleep(1000);
                    observable.OnNext("b");
                    Thread.Sleep(1000);
                    observable.OnCompleted();
                    Thread.Sleep(1000);
                });

                return Disposable.Create(() => Console.WriteLine("Observer has unsubscribed"));
            });
    }

}

订阅者

public static class Subscriber{
public static bool Subscribe()
    {
        Publisher.NonBlocking().Subscribe((s) =>
        {
            Debug.WriteLine(s);
        }, () =>
        {
            Debug.WriteLine("Complete");
        });
        // This will currently return true before the subscription is complete
        // I want to block and not Return until the Subscriber is Complete
        return true;
    }

}

【问题讨论】:

    标签: .net system.reactive reactive-programming


    【解决方案1】:

    您需要为此使用System.Reactive.Threading.Task

    把你的 observable 变成一个任务...

    var source = Publisher.NonBlocking()
        .Do(
            (s) => Debug.WriteLines(x),
            ()  => Debug.WriteLine("Completed")
        )
        .LastOrDefault()
        .ToTask();
    

    Do(...).Subscribe() 就像Subscribe(...)。所以Do 只是增加了一些副作用。

    LastOrDefault 在那里是因为ToTask 创建的Task 将只等待来自源Observable 的第一项,如果没有生成任何项,它将失败(抛出)。因此,LastOrDefault 有效地导致 Task 等待源完成,无论它产生什么。

    所以我们有了任务之后,就等着吧:

    task.Wait(); // blocking
    

    或者使用异步/等待:

    await task; // non-blocking
    

    编辑:

    Cory Nelson 提出了一个很好的观点:

    在最新版本的 C# 和 Visual Studio 中,您实际上可以awaitIObservable&lt;T&gt;。这是一个非常很酷的功能,但它的工作方式与等待Task 略有不同。

    当您等待任务时,它会导致任务运行。如果多次等待某个任务的单个实例,则该任务将只执行一次。 Observables 略有不同。您可以将可观察对象视为具有多个返回值的异步函数……每次订阅可观察对象时,可观察对象/函数都会执行。所以这两段代码有不同的含义:

    等待 Observable:

    // Console.WriteLine will be invoked twice.
    var source = Observable.Return(0).Do(Console.WriteLine);
    await source; // Subscribe
    await source; // Subscribe
    

    通过任务等待 Observable:

    // Console.WriteLine will be invoked once.
    var source = Observable.Return(0).Do(Console.WriteLine);
    var task = source.ToTask();
    await task; // Subscribe
    await task; // Just yield the task's result.
    

    因此,从本质上讲,等待 Observable 的工作方式如下:

    // Console.WriteLine will be invoked twice.
    var source = Observable.Return(0).Do(Console.WriteLine);
    await source.ToTask(); // Subscribe
    await source.ToTask(); // Subscribe
    

    但是,await observable 语法在 Xamerin Studio 中不起作用(截至撰写本文时)。如果您使用的是 Xamerin Studio,我强烈建议您在最后一刻使用 ToTask,以模拟 Visual Studio 的 await observable 语法的行为。

    【讨论】:

    • 太棒了,不知道 '.Do(..)' 看起来也存在 LastOrDefaultAsync() 方法,所以我可以从中执行 .Wait()
    • 您可以直接awaitIObservable&lt;&gt;。它将返回序列中的最后一项。
    • @Lukie,在使用 LastOrDefaultAsync 之前,您可以看看我的编辑。过早将 Observable 转换为 Task 会对代码的可重用性产生深远影响。话虽如此,如果您在最后一刻使用 LastOrDefaultAsync,那可能是您最好的选择。
    猜你喜欢
    • 1970-01-01
    • 2023-03-04
    • 1970-01-01
    • 2013-05-15
    • 2011-11-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多