【问题标题】:How do I create an Rx sequence by running tasks over original sequence's values?如何通过在原始序列的值上运行任务来创建 Rx 序列?
【发布时间】:2013-07-11 02:23:50
【问题描述】:

我有一个IObservable<T> 类型的序列和一个将T, CancellationToken 映射到Task<U> 的函数。从他们那里获得IObservable<U> 的最干净的方法是什么?

我需要以下语义:

  • 每个任务在前一个项目的任务完成后开始
  • 如果任务已被取消或出错,则会被跳过
  • 严格保留原始序列的顺序

这是我看到的签名:

public static IObservable<U> Select<T, U> (
    this IObservable<T> source,
    Func<T, CancellationToken, Task<U>> selector
);

我还没有写任何代码,但除非有人打败我,否则我会写。
无论如何,我不熟悉像Window 这样的运算符,所以我的解决方案可能不太优雅。

我需要 C# 4 中的解决方案,但为了比较,也欢迎 C# 5 答案。


如果你好奇,下面是我的真实场景,或多或少:

Dropbox.GetImagesRecursively ()
    .ObserveOn (SynchronizationContext.Current)
    .Select (DownloadImage)
    .Subscribe (AddImageToFilePicker);

【问题讨论】:

  • 哇,我刚刚意识到SelectMany 具有完全相同的签名。唉,它不会等待任务完成。
  • 我认为一般来说,忽略异常是个坏主意。他们至少应该被记录下来,或者类似的东西。
  • @svick:你说得对,我会在第一次实施后修改这个设计。

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


【解决方案1】:

到目前为止,这似乎对我有用:

public static IObservable<U> Select<T, U> (
    this IObservable<T> source,
    Func<T, CancellationToken, Task<U>> selector)
{
    return source
        .Select (item => 
            Observable.Defer (() => 
                Observable.StartAsync (ct => selector (item, ct))
                    .Catch (Observable.Empty<U> ())
            ))
        .Concat ();
}

我们将一个基于延迟任务的异常吞咽可观察对象映射到每个项目,然后将它们连接起来。


我的思考过程是这样的。

我注意到其中一个SelectMany 重载几乎完全符合我的要求,甚至具有完全相同的签名。但它并没有满足我的需求:

  • 它会在原始项目出现时创建任务,而我需要等待每个任务完成
  • 它不提供跳过已取消和错误任务的选项

我查看了这个重载的实现,发现它使用FromAsync 来处理任务创建和取消:

public virtual IObservable<TResult> SelectMany<TSource, TTaskResult, TResult> (IObservable<TSource> source, Func<TSource, CancellationToken, Task<TTaskResult>> taskSelector, Func<TSource, TTaskResult, TResult> resultSelector)
{
    return SelectMany_<TSource, TTaskResult, TResult> (
        source,
        x => FromAsync (ct => taskSelector (x, ct)),
        resultSelector
    );
}

我将目光转向FromAsync 以了解它是如何实现的,并且惊喜地发现它也是可组合的:

public virtual IObservable<TResult> FromAsync<TResult> (Func<CancellationToken, Task<TResult>> functionAsync)
{
    return Defer (() => StartAsync (functionAsync));
}

我重用了DeferStartAsync,同时还添加了Catch 以吞下错误。 DeferConcat 的组合确保任务相互等待并按原始顺序启动。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2022-12-19
    • 1970-01-01
    • 2012-05-30
    • 1970-01-01
    • 2019-11-20
    • 1970-01-01
    • 2014-01-28
    • 1970-01-01
    相关资源
    最近更新 更多