【问题标题】:'WaitFor' an observable'等待'一个可观察的
【发布时间】:2012-07-27 21:08:51
【问题描述】:

我有一个正在处理的任务列表(启用驱动器、更改位置、等待停止、禁用)。

“等待”监控我想要等待的IObservable<Status>(这样我就可以将它通过ContinueWith 和其他任务线程化)。

我开始在订阅者的 OnNext 处理内执行以下任务,但这很丑陋。我现在想出的是这个扩展方法:

public static Task<T> WaitFor<T>(this IObservable<T> source, Func<T, bool> pred)
{
    var tcs = new TaskCompletionSource<T>();
    source
        .Where(pred)
        .DistinctUntilChanged()
        .Take(1)  //OnCompletes the observable, subscription will self-dispose
        .Subscribe(val => tcs.TrySetResult(val),
                    ex => tcs.TrySetException(ex),
                    () => tcs.TrySetCanceled());

    return tcs.Task;
}

(已更新 svick 建议处理OnCompleted 和OnError)

问题:

  • 这是好事、坏事还是丑陋?
  • 我是否错过了可以做到这一点的现有扩展?
  • Where 和 DistinctUntilChanged 的顺序是否正确? (我认为他们是)

【问题讨论】:

  • 你不应该也处理source出错或完成的情况吗?
  • @svick,嗯,好问题。在我的特定用例中,我正在观察Repeat,所以我认为这不会完成。但确实,如果我希望这个扩展是可重用的,我应该处理这些情况。

标签: c# task-parallel-library system.reactive


【解决方案1】:

对此不是 100% 确定,但通过阅读 Rx 2.0 beta 博客文章,我认为如果您可以使用 async/await,您可以“返回 await source.FirstAsync(pred)”或不使用 async,“return source.FirstAsync(pred).ToTask()"

http://blogs.msdn.com/b/rxteam/archive/2012/03/12/reactive-extensions-v2-0-beta-available-now.aspx

【讨论】:

  • 当/如果我升级时,我会记住这一点。目前,我正在尝试查看 没有 async/await .. 有多难
【解决方案2】:

至少我会把这个扩展方法改成这样:

public static Task<T> WaitFor<T>(this IObservable<T> source, Func<T, bool> pred)
{
    return
        source
            .Where(pred)
            .DistinctUntilChanged()
            .Take(1)
            .ToTask();
}

使用.ToTask() 比引入TaskCompletionSource 要好得多。您需要引用 System.Reactive.Threading.Tasks 命名空间来获取 .ToTask() 扩展方法。

此外,DistinctUntilChanged 在此代码中是多余的。您只能获得一个值,因此默认情况下它必须是不同的。

现在,我的下一个建议可能有点争议。这个扩展是个坏主意,因为它隐藏了正在发生的事情的真正语义。

如果我有这两个代码片段:

var t = xs.WaitFor(x => x > 10);

或者:

var t = xs.Where(x => x > 10).Take(1).ToTask();

我通常更喜欢第二个片段,因为它清楚地向我展示了正在发生的事情 - 我不需要记住 WaitFor 的语义。

除非您使 WaitFor 的名称更具描述性 - 也许是 TakeOneAsTaskWhere - 那么您将在使用它的代码中明确使用运算符并使代码更难管理。

以下内容不是更容易记住语义吗?

var t = xs.TakeOneAsTaskWhere(x => x > 10);

对我来说,底线是 Rx 运算符是要组合的,而不是封装的,但如果你要封装它们,那么它们的含义必须清楚。

我希望这会有所帮助。

【讨论】:

  • 谢谢。我在网络上的不同地方看到过ToTask,但(懒惰地)认为它在过去或以前的版本中丢失了。我认为你是对的,使用ToTask 代码变得不那么难看,因此减少了对扩展方法的需求。
  • 我认为封装可以很好地应用一组组合在一起时具有更高含义的组合。但我认为这里没有足够的“更高含义”值得封装 - 封装的名称必须比一组组合更有意义。如果将 xs.Where(...).Take(1).ToTask() 替换为 xs.FirstAsync(...).ToTask(),那将非常简洁明了。
猜你喜欢
  • 2012-05-04
  • 2017-09-08
  • 2022-08-19
  • 1970-01-01
  • 1970-01-01
  • 2018-01-03
  • 1970-01-01
  • 2019-03-25
  • 2021-06-10
相关资源
最近更新 更多