【问题标题】:Using Rx create a polling request for webservice call使用 Rx 创建 web 服务调用的轮询请求
【发布时间】:2016-02-09 07:41:11
【问题描述】:

在 C# 中使用 Rx 我正在尝试创建对 REST API 的轮询请求。我面临的问题是,Observable 需要按顺序发送响应。表示如果请求 A 在 X 时间到达,请求 B 在 X + dx 时间到达,并且 B 的响应在 A 之前到达,则 Observable 表达式应忽略或取消请求 A。

我编写了一个示例代码来尝试描述该场景。如何修复它以仅获取最新响应并取消或忽略以前的响应。

 class Program
    {
        static int i = 0;

        static void Main(string[] args)
        {
            GenerateObservableSequence();

            Console.ReadLine();
        }

        private static void GenerateObservableSequence()
        {
            var timerData = Observable.Timer(TimeSpan.Zero,
                TimeSpan.FromSeconds(1));

            var asyncCall = Observable.FromAsync<int>(() =>
            {
                TaskCompletionSource<int> t = new TaskCompletionSource<int>();
                i++;

                int k = i;
                var rndNo = new Random().Next(3, 10);
                Task.Delay(TimeSpan.FromSeconds(rndNo)).ContinueWith(r => { t.SetResult(k); });
                return t.Task;
            });

            var obs = from t in timerData
            from data in asyncCall
            select data;

            var hot = obs.Publish();
            hot.Connect();

                hot.Subscribe(j => 
            {
                Console.WriteLine("{0}", j);
            });
        }
    }

@Enigmativity 回答后:添加 Polling Aync 功能以始终获取最新响应:

 public static IObservable<T> PollingAync<T> (Func<Task<T>> AsyncCall, double TimerDuration)
        {
            return Observable
         .Create<T>(o =>
         {
             var z = 0L;
             return
                 Observable
                     .Timer(TimeSpan.Zero, TimeSpan.FromSeconds(TimerDuration))
                     .SelectMany(nr =>
                         Observable.FromAsync<T>(AsyncCall),
                         (nr, obj) => new { nr, obj})
                     .Do(res => z = Math.Max(z, res.nr))
                     .Where(res => res.nr >= z)
                     .Select(res => res.obj)
                     .Subscribe(o);
         });

    }

【问题讨论】:

    标签: c# recursion system.reactive reactive-programming polling


    【解决方案1】:

    这是一种常见的情况,可以简单地修复。

    您的示例代码的关键部分是

    var obs = from t in timerData
              from data in asyncCall
              select data;
    

    这可以理解为“对于timerData 中的每个值,获取asyncCall 中的所有值”。这是SelectMany(或FlatMap)运算符。 SelectMany 运算符将从内部序列 (asyncCall) 中获取所有值,并在收到时返回它们的值。这意味着您可以得到乱序值。

    当外部序列 (timerData) 产生新值时,您想要取消之前的内部序列。为此,我们想改用Switch 运算符。

    var obs = timerData.Select(_=>asyncCall)
                       .Switch();
    

    完整的代码可以清理为以下内容。 (删除了多余的发布/连接,按键处理订阅)

    类程序 { 静态 int i = 0;

        static void Main(string[] args)
        {
            using (GenerateObservableSequence().Subscribe(x => Console.WriteLine(x)))
            {
                Console.ReadLine();
            }
        }
    
        private static IObservable<int> GenerateObservableSequence()
        {
            var timerData = Observable.Timer(TimeSpan.Zero,
                TimeSpan.FromSeconds(1));
    
            var asyncCall = Observable.FromAsync<int>(() =>
            {
                TaskCompletionSource<int> t = new TaskCompletionSource<int>();
                i++;
    
                int k = i;
                var rndNo = new Random().Next(3, 10);
                Task.Delay(TimeSpan.FromSeconds(rndNo)).ContinueWith(r => { t.SetResult(k); });
                return t.Task;
            });
    
            return from t in timerData
                   from data in asyncCall
                   select data;
        }
    }
    

    --编辑--

    看来我误解了这个问题。 @Enigmativity 提供了更准确的答案。这是对他答案的清理。

    //Probably should be a field?
    var rnd = new Random();
    var obs = Observable.Timer(TimeSpan.Zero, TimeSpan.FromSeconds(1))
            //.Select(n => new { n, r = ++i })
            //No need for the `i` counter. Rx does this for us with this overload of `Select`
            .Select((val, idx) => new { Value = val, Index = idx})
            .SelectMany(nr =>
                Observable.Timer(TimeSpan.FromSeconds(rnd.Next(3, 10))),
                (nr, _) => nr)
            //.Do(nr => z = Math.Max(z, nr.n))
            //.Where(nr => nr.n >= z)
            //Replace external State and Do with scan and Distinct
            .Scan(new { Value = 0L, Index = -1 }, (prev, cur) => {
                return cur.Index > prev.Index
                    ? cur
                    : prev;
            })
            .DistinctUntilChanged()
            .Select(nr => nr.Value)
            .Dump();
    

    【讨论】:

    • 这不是总是只取最新请求的值吗?我对这个问题的理解是,如果发生两次对asyncCall 的调用并且结果按调用顺序返回,那么这两个结果都是有效的,但是如果第二次调用首先返回,则忽略第一次调用的结果。跨度>
    • 根据您对问题的理解,那么您的假设也是正确的。
    • @LeeCampbell 关于如何处理我们忽略的先前调用的异常的任何建议。而不是 Observable.Timer(TimeSpan.FromSeconds(rnd.Next(3, 10)) 我将调用一个异步函数,如更新的问题所示。这个异步函数可能会引发异常。因此我希望取消之前的异步调用并处理异常而不结束 IObservable 流。
    • 一般来说,Rx 是关于组合序列的,一旦序列出错,它就被认为是终端。因此,从这个意义上说,您的用例不受支持。但是,在这种情况下,您可以使用 Catch 运算符并从那里产生适当的值。见introtorx.com/Content/v1.0.10621.0/…
    【解决方案2】:

    让我们从简化代码开始。

    这基本上是相同的代码:

    var rnd = new Random();
    
    var i = 0;
    
    var obs =
        from n in Observable.Timer(TimeSpan.Zero, TimeSpan.FromSeconds(1))
        let r = ++i
        from t in Observable.Timer(TimeSpan.FromSeconds(rnd.Next(3, 10)))
        select r;
    
    obs.Subscribe(Console.WriteLine);
    

    我得到这样的结果:

    2 1 3 4 8 5 11 6 9 7 10

    或者,这可以写成:

    var obs =
        Observable
            .Timer(TimeSpan.Zero, TimeSpan.FromSeconds(1))
            .Select(n => ++i)
            .SelectMany(n =>
                Observable.Timer(TimeSpan.FromSeconds(rnd.Next(3, 10))), (n, _) => n);
    

    所以,现在满足您的要求:

    如果请求 A 在 X 时间到达并且请求 B 在 X + dx 时间到达并且 B 的响应在 A 之前到达,则 Observable 表达式应该忽略或取消请求 A。

    代码如下:

    var rnd = new Random();
    
    var i = 0;
    var z = 0L;
    
    var obs =
        Observable
            .Timer(TimeSpan.Zero, TimeSpan.FromSeconds(1))
            .Select(n => new { n, r = ++i })
            .SelectMany(nr =>
                Observable.Timer(TimeSpan.FromSeconds(rnd.Next(3, 10))), (nr, _) => nr)
            .Do(nr => z = Math.Max(z, nr.n))
            .Where(nr => nr.n >= z)
            .Select(nr => nr.r);
    

    我不喜欢这样使用.Do,但我还想不出替代方案。

    这给出了这种东西:

    1 5 8 9 10 11 14 15 16 17 22

    请注意,这些值只是升序。

    现在,您确实应该使用Observable.Create 来封装您正在使用的状态。所以你最终的 observable 应该是这样的:

    var obs =
        Observable
            .Create<int>(o =>
            {
                var rnd = new Random();
                var i = 0;
                var z = 0L;
                return
                    Observable
                        .Timer(TimeSpan.Zero, TimeSpan.FromSeconds(1))
                        .Select(n => new { n, r = ++i })
                        .SelectMany(nr =>
                            Observable.Timer(TimeSpan.FromSeconds(rnd.Next(3, 10))),
                            (nr, _) => nr)
                        .Do(nr => z = Math.Max(z, nr.n))
                        .Where(nr => nr.n >= z)
                        .Select(nr => nr.r)
                        .Subscribe(o);
            });
    

    【讨论】:

    • 我们如何避免来自使用 where 子句过滤的旧响应的异常?
    • @BalrajSingh - 那么,如果我用偶尔会引发错误的操作替换++i?这是可能的情况吗?会发生什么样的错误?
    • 我的意思是代替 Observable.Timer(TimeSpan.FromSeconds(rnd.Next(3, 10)) 我将用对 Web 服务的异步调用替换它,它可能会引发任何类型的 http 错误。在这种情况下,我需要取消我发出的所有旧请求,这样它就不会出现任何错误,如果它抛出,那么我应该有办法在不结束我的 Observable 流的情况下处理它。
    • 我已经用示例代码更新了问题,用于调用任何异步函数并仅获取最新响应。如果异步函数这样做,这可能会引发异常。
    • Enigmativity,我提供了一个没有外部变量和Do 语句的答案的实现。开箱即用,做得很好。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-11-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多