【问题标题】:Throttle IObservable based on whether async handler is still busy根据异步处理程序是否仍然忙来限制 IObservable
【发布时间】:2015-07-23 02:47:32
【问题描述】:

我有一个IObservable,它每秒生成一个值,然后是一个运行可能需要一些时间的代码的选择:

var events = Observable.Interval(TimeSpan.FromSeconds(1));
ssoInfoObservable = events
    .Select(async e =>
    {
        Console.Out.WriteLine("Select   : " + e);
        await Task.Delay(4000);
        return e;
    })
    .SelectMany(t => t.ToObservable())
    .Subscribe(l => Console.WriteLine("Subscribe: " + l));

在我的示例中,长时间运行的操作需要 4 秒。当Select 中的代码正在运行时,我不希望从Interval 中生成另一个值。我该如何做到这一点?这可能吗?也许使用特定的IScheduler 实现?

请注意,如果没有异步代码,一切都会按照here 所述的预期工作。

这个问题和one I asked earlier非常相似,除了async/await

【问题讨论】:

  • .ObserveOn(...) 没有任何区别。 .Interval(...) 运算符计算每组处理程序之间的间隔(即间隔),而不是每个开始时间之间的间隔。
  • 我发现我的示例代码太简单了,我更新了我的问题以反映我的实际代码。
  • 这仍然受限于实际运行的线程数。我每个人都只有四个选择,然后就一直呆在那里。我会考虑一下,但你能描述一下为什么这是你的实际代码吗?您认为它可以为您解决什么问题?
  • 最重要的要求是Select 不能并行运行。等待的操作可能会长时间运行(等待结果在数据库中可用)并为链的其余部分生成结果。简单的解决方案当然是同步等待等待任务的Result。实际代码非常相似,只是我在等待一些有用的东西。
  • 你为什么要像await 这样混合 Observables? Observables 实现了将您的工作推送到后台线程的目的,从而释放 UI(或主线程)以响应其他代码。混入await 有点像“双重浸渍”。

标签: c# async-await system.reactive throttling


【解决方案1】:

请参阅此示例以创建 async generate function。你会略有不同,因为你需要一个时间偏移,你不需要 iterate 所以它看起来更像这样:

    public static IObservable<T> GenerateAsync<T>(TimeSpan span, 
                                                  Func<int, Task<T>> generator, 
                                                  IScheduler scheduler = null)
    {
        scheduler = scheduler ?? Scheduler.Default;

        return Observable.Create<T>(obs =>
        {
            return scheduler.Schedule(0, span, async (idx, recurse) =>
            {
                obs.OnNext(await generator(idx));
                recurse(idx + 1, span);
            });

        });
    }

用法:

  Extensions.GenerateAsync(TimeSpan.FromSeconds(1), idx => /*Async work, return a task*/, scheduler);

作为可能的第二种选择,您可以将switchFirst 的实现移植到C# 中。 SwitchFirst 将订阅它收到的第一个 Observable 并忽略后续的,直到当前订阅完成。

如果你采用这种方法,你可能会得到类似的东西:

Observable.Interval(TimeSpan.FromSeconds(1))
          .Select(e => Observable.FromAsync(() => /*Do Async stuff*/)
          .SwitchFirst()
          .Subscribe();

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-11-28
    • 1970-01-01
    • 1970-01-01
    • 2017-07-28
    • 2011-06-23
    • 1970-01-01
    相关资源
    最近更新 更多