【问题标题】:How do you chain together two Asynchronous Operations with the Reactive Framework?您如何使用反应式框架将两个异步操作链接在一起?
【发布时间】:2010-06-14 20:24:54
【问题描述】:

我真正想做的就是完成两个异步操作,一个接一个。例如

下载网站 X。完成后,下载网站 Y。

【问题讨论】:

    标签: .net asynchronous system.reactive


    【解决方案1】:

    从两个 observable 中执行 SelectMany(使用 ToAsync 将操作作为 IObservable 获取),就像 LINQ 选择许多(或使用语法糖):

    var result = 
    from x in X
    from y in Y
    select new { x, y }  
    

    还有其他选择,但这取决于您的具体情况。

    【讨论】:

    • 请注意。如果“X”或“Y”多次返回,我们将开始看到交叉连接行为。它非常适合只返回一次的异步操作。但为了安全起见,使用 Prune 扩展方法将只允许通过一个通知,因此对于诸如 web 服务调用之类的东西是一个好主意。
    【解决方案2】:

    大概是这样的:

        static IObservable<DownloadProgressChangedEventArgs> CreateDownloadFileObservable(string url, string fileName)
        {
           IObservable<DownloadProgressChangedEventArgs> observable =
               Observable.CreateWithDisposable<DownloadProgressChangedEventArgs>(o =>
            {
                var cancellationTokenSource = new CancellationTokenSource();
                Scheduler.TaskPool.Schedule
                (
                    () =>
                    {
    
                        Thread.Sleep(3000);
                        if (!cancellationTokenSource.Token.IsCancellationRequested)
                        {
                            WebClient client = new WebClient();
                            client.DownloadFileAsync(new Uri(url), fileName,fileName);
    
    
                            DownloadProgressChangedEventHandler prgChangedHandler = null;
                            prgChangedHandler = (s,e) =>
                                {                                                   
                                    o.OnNext(e);
                                };
    
    
                            AsyncCompletedEventHandler handler = null;
                            handler = (s, e) =>
                            {
                                prgChangedHandler -= prgChangedHandler;                            
                                if (e.Error != null)
                                {
                                    o.OnError(e.Error);
                                }
                                client.DownloadFileCompleted -= handler;
                                o.OnCompleted();
                            };
                            client.DownloadFileCompleted += handler;
                            client.DownloadProgressChanged += prgChangedHandler;
                        }
                        else
                        {
                            Console.WriteLine("Cancelling download of {0}",fileName);
                        }
                    }
                );
                return cancellationTokenSource;
            }
            );
           return observable;
        }
    
        static void Main(string[] args)
        {
    
            var obs1 = CreateDownloadFileObservable("http://www.cnn.com", "cnn.htm");
            var obs2 = CreateDownloadFileObservable("http://www.bing.com", "bing.htm");
    
            var result = obs1.Concat(obs2);
            var subscription = result.Subscribe(a => Console.WriteLine("{0} -- {1}% complete ",a.UserState,a.ProgressPercentage), e=> Console.WriteLine(e.Message),()=> Console.WriteLine("Completed"));
            Console.ReadKey();
            subscription.Dispose();
            Console.WriteLine("Press a key to exit");
            Console.ReadKey();
    }
    

    你可以把 CreateDownloadFileObservable 变成 WebClient 上的扩展方法,去掉方法定义中的WebClient client = new WebClient();,客户端会变成这样:

    WebClient c = new WebClient();
    c.CreateDownloadFileObservable ("www.bing.com","bing.htm");
    

    【讨论】:

      【解决方案3】:

      我自己不喜欢这样的建议,但是...

      如果您需要按顺序执行两个操作,请按顺序执行(您仍然可以在不同于 main 的线程上执行它们)。

      如果您仍想在代码中分离任务,那么来自 System.Parallel 的新结构将是合适的:

      var task1 = Task.Factory.StartNew (() => FirstTask());
      var task2 = task1.ContinueWith (frst => SecondTask ());
      

      如果这是一种了解 system.reactive 的方法,请尝试 Observable.GenerateInSequence - 但它肯定会过大。请记住,observable 是 enumerable 的对应物,最好以类似的方式使用(迭代数据)。它不是运行异步操作的工具。

      编辑:我想承认我错了,Richard 是对的——在我回复的时候,我对 RX 并不完全满意。现在我认为RX是启动异步操作最自然的方式。不过,Richard 的回答有点吝啬,应该是:

      var result =  
      from x in XDownload.ToAsync()() 
      from y in YDownload.ToAsync()() 
      select y   
      

      【讨论】:

      • "它不是运行异步操作的工具。"我完全不同意这种说法,延续单子(Rx)正是我们组合异步操作所需要的,因为我们显式地传递任何状态并且所有副作用都得到管理。首先创建 Rx 的主要原因之一是异步和并行编程。
      • 观察异步操作 - 是的。从可观察管道创建异步操作 - 是的。使用 Observable 启动异步任务(尤其是一两个)而不观察结果 - 这将过于复杂。
      • 是的,Rx 的重点是异步操作的可组合性。我刚写完一个严重依赖 Rx 的应用程序,我使用的代码看起来与 Richards 的示例一模一样。在一个地方,我需要发出多达 18 个 Web 服务请求来加载数据,如果没有 Rx,我需要 10 倍的代码。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2018-06-04
      • 1970-01-01
      • 1970-01-01
      • 2017-09-26
      • 1970-01-01
      • 2021-08-01
      • 2018-11-23
      相关资源
      最近更新 更多