【问题标题】:How to manage observable subscription for dependent observables?如何管理依赖 observable 的 observable 订阅?
【发布时间】:2014-08-22 19:26:25
【问题描述】:

这个示例控制台应用程序有 2 个 observables。第一个将数字从 1 推送到 100。这个 observable 由AsyncClass 订阅,它为它获得的每个数字运行一个长时间运行的过程。在完成这个新的异步过程后,我希望能够“推送”给 2 个订阅者,他们将使用这个新值做一些事情。

下面的源代码中注释了我的尝试。

异步类:

class AsyncClass
    {
        private readonly IConnectableObservable<int> _source;
        private readonly IDisposable _sourceDisposeObj;
        public IObservable<string> _asyncOpObservable;

        public AsyncClass(IConnectableObservable<int> source)
        {
            _source           = source;
            _sourceDisposeObj = _source.Subscribe(
                                    ProcessArguments,
                                    ExceptionHandler,
                                    Completed
                                );
            _source.Connect();


        }

        private void Completed()
        {
            Console.WriteLine("Completed");
            Console.ReadKey();
        }

        private void ExceptionHandler(Exception exp)
        {
            throw exp;
        }

        private void ProcessArguments(int evtArgs)
        {
            Console.WriteLine("Argument being processed with value: " + evtArgs);

            //_asyncOpObservable = LongRunningOperationAsync("hello").Publish(); 
            // not going to work either since this creates a new observable for each value from main observer

        }
        // http://rxwiki.wikidot.com/101samples
        public IObservable<string> LongRunningOperationAsync(string param)
        {
            // should not be creating an observable here, rather 'pushing' values?
            return Observable.Create<string>(
                o => Observable.ToAsync<string, string>(DoLongRunningOperation)(param).Subscribe(o)
            );
        }

        private string DoLongRunningOperation(string arg)
        {
            return "Hello";
        }
    }

主要:

static void Main(string[] args)
        {
            var source = Observable
                            .Range(1, 100)
                            .Publish();

            var asyncObj = new AsyncClass(source);
            var _asyncTaskSource = asyncObj._asyncOpObservable;

            var ui1 = new UI1(_asyncTaskSource);
            var ui2 = new UI2(_asyncTaskSource);
        }

UI1(和UI2,基本一样):

 class UI1
    {
        private IConnectableObservable<string> _asyncTaskSource;
        private IDisposable _taskSourceDisposable;

        public UI1(IConnectableObservable<string> asyncTaskSource)
        {
            _asyncTaskSource = asyncTaskSource;
            _asyncTaskSource.Connect();
            _taskSourceDisposable = _asyncTaskSource.Subscribe(RefreshUI, HandleException, Completed);
        }

        private void Completed()
        {
            Console.WriteLine("UI1: Stream completed");
        }

        private void HandleException(Exception obj)
        {
            Console.WriteLine("Exception! "+obj.Message);
        }

        private void RefreshUI(string obj)
        {
            Console.WriteLine("UI1: UI refreshing with value "+obj);
        }
    }

这是我第一个使用 Rx 的项目,所以如果我应该换个思路,请告诉我。任何帮助将不胜感激!

【问题讨论】:

    标签: c# system.reactive reactive-programming


    【解决方案1】:

    我要让你知道你应该有不同的想法...... :) 抛开轻率不谈,这看起来像是面向对象和函数响应式风格之间的严重冲突。

    这里不清楚数据流时序和结果缓存的要求是什么——Publish 和IConnectableObservable 的使用有点混乱。我猜你想避免 2 个下游订阅导致处理重复的值?我的一些答案是基于这个前提。使用Publish() 可以通过允许多个订阅者共享对单个源的订阅来实现这一点。

    Idiomatic Rx 希望您尝试并保持实用风格。为此,您希望将长期运行的工作呈现为一个函数。因此,我们可以说,与其尝试将 AsyncClass 逻辑作为一个类直接连接到 Rx 链中,不如将其呈现为一个函数,就像这个人为的示例一样:

    async Task<int> ProcessArgument(int argument)
    {
        // perform your lengthy calculation - maybe in an OO style,
        // maybe creating class instances and invoking methods etc.
        await Task.Delay(TimeSpan.FromSeconds(1));
        return argument + 1;
    }
    

    现在,您可以构造一个完整的 Rx 可观察链来调用该函数,并且通过使用 Publish().RefCount() 可以避免多个订阅者造成重复工作。注意这也是如何分离关注点的——处理值的代码更简单,因为重用是在其他地方处理的。

    var query = source.SelectMany(x => ProcessArgument(x).ToObservable())
                      .Publish().RefCount();
    

    通过为订阅者创建单个链,仅在订阅时需要时才开始工作。我使用了Publish().RefCount() - 但如果你想确保第二个和后续订阅者不会错过值,你可以使用Replay(简单)或使用Publish() 然后Connect - 但你会想要在单个订阅者代码之外的 Connect 逻辑,因为您只需要在所有订阅者都订阅后调用它一次。

    【讨论】:

    • 感谢您的回答詹姆斯,它帮助我了解我做错了什么。如果您仍然想知道,Observable.Range 只是部分连续事件(如“文件保存”)的虚拟可观察对象。当我收到任何测试文件的文件保存事件时,我正在尝试缩小一些测试文件。之后,它应该通知 2 个 UI 类该文件已与文件内容一起保存。我展示的代码是一个人为的例子。
    • 刚刚注意到SelectMany 在将其发送到管道之前会累积所有值。我想在数据到达时发送它。有什么建议吗?
    • SelectMany 肯定不会这样做,肯定是其他人负责...很难说没有看到您的代码,但请检查您的假设。
    • 这可能有助于调试! stackoverflow.com/questions/20220755/…
    • 我现在觉得自己很蠢,你说得对。调试脚本很棒。感谢您在这里发布:)
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2013-12-03
    • 1970-01-01
    • 2023-04-02
    • 2018-08-28
    • 2019-12-29
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多