【问题标题】:What is the Rx.NET way to produce a cancellable observable of file names?Rx.NET 生成可取消的可观察文件名的方法是什么?
【发布时间】:2015-07-20 17:33:38
【问题描述】:

我想生成一个可观察的文件,以便可以随时取消文件名的发现。在本示例中,取消会在 1 秒内自动发生。

这是我当前的代码:

class Program
{
    static void Main()
    {
        try
        {
            RunAsync(@"\\abc\xyz").GetAwaiter().GetResult();
        }
        catch (Exception exc)
        {
            Console.Error.WriteLine(exc);
        }
        Console.Write("Press Enter to exit");
        Console.ReadLine();
    }

    private static async Task RunAsync(string path)
    {
        var cts = new CancellationTokenSource(TimeSpan.FromSeconds(1));
        await GetFileSource(path, cts);
    }

    private static IObservable<string> GetFileSource(string path, CancellationTokenSource cts)
    {
        return Observable.Create<string>(obs => Task.Run(async () =>
        {
            Console.WriteLine("Inside Before");
            foreach (var file in Directory.EnumerateFiles(path, "*", SearchOption.AllDirectories).Take(50))
            {
                cts.Token.ThrowIfCancellationRequested();
                obs.OnNext(file);
                await Task.Delay(100);
            }
            Console.WriteLine("Inside After");
            obs.OnCompleted();
            return Disposable.Empty;
        }, cts.Token))
        .Do(Console.WriteLine);
    }
}

我不喜欢我的实现的两个方面(如果有更多 - 请随时指出):

  1. 我有一个可枚举的文件,但我手动迭代每个文件。我可以以某种方式使用ToObservable 扩展吗?
  2. 我不知道如何使用传递给Task.Runcts.Token。必须使用从外部上下文捕获的ctsGetFileSource 参数)。我觉得很难看。

这是应该怎么做的吗?一定是更好的方法。

【问题讨论】:

  • 这似乎不是一个非常被动的问题,因为您实际上只是在枚举集合。是什么导致取消?您是否看过 Parallel.ForEach 或 PLinq,它们也支持中间迭代取消?
  • 这是一个精简的例子。真正的逻辑要复杂得多。
  • 作为一般规则 - 如果您发现自己在做return Disposable.Empty;,那么您几乎可以肯定做错了什么。

标签: c# system.reactive


【解决方案1】:

我仍然不相信这真的是一个反应性问题,你要求对生产者施加背压,这确实违背了反应性应该如何工作。

话虽如此,如果你打算这样做,你应该意识到非常细粒度的时间操作应该几乎总是委派给Scheduler,而不是尝试与TasksCancellationTokens进行协调.所以我会重构成这样:

public static IObservable<string> GetFileSource(string path, Func<string, Task<string>> processor, IScheduler scheduler = null) {

  scheduler = scheduler ?? Scheduler.Default;

  return Observable.Create<string>(obs => 
  {
    //Grab the enumerator as our iteration state.
    var enumerator = Directory.EnumerateFiles(path, "*", SearchOption.AllDirectories)
                              .GetEnumerator();
    return scheduler.Schedule(enumerator, async (e, recurse) =>
    {
      if (!e.MoveNext())
      {
         obs.OnCompleted();
         return;
      }

      //Wait here until processing is done before moving on
      obs.OnNext(await processor(e.Current));

      //Recursively schedule
      recurse(e);
    });
  });

}

然后,不要传入取消令牌,而是使用TakeUntil

var source = GetFileSource(path, x => {/*Do some async task here*/; return x; })
 .TakeUntil(Observable.Timer(TimeSpan.FromSeconds(1));

您还可以查看更高级的示例,以实现 async Generate method

【讨论】:

  • 哇。这是一个很棒的方法。
【解决方案2】:

当您可以使用其他运算符时,我建议您避免使用Observable.Create

此外,当您在 Observable.Create 中执行 return Disposable.Empty; 时,您正在创建一个无法被普通 Rx 订阅一次性停止的 observable。这可能会导致内存泄漏和不必要的处理。

最后,抛出异常来结束正常计算是个坏主意。

有一个很好的干净解决方案似乎可以满足您的需求:

private static IObservable<string> GetFileSource(string path, CancellationTokenSource cts)
{
    return
        Directory
            .EnumerateFiles(path, "*", SearchOption.AllDirectories)
            .ToObservable()
            .Take(50)
            .TakeWhile(f => !cts.IsCancellationRequested);
}

我唯一没有包括的是Task.Delay(100);。你为什么要这么做?

【讨论】:

  • 延迟是为了模拟每个获取的文件发生的异步工作。
  • 您的解决方案的语义略有不同,尚不确定它是否适合我。在原始示例中,文件枚举器不会继续到下一个文件,直到某些异步处理结束(使用 Task.Delay 模拟)。在你的代码中(请删除Take(50),因为这个数字没有任何意义)我可以在TakeWhile之后添加SelectMany来运行异步后处理,但它不会一样,因为文件枚举器是将获取更多文件,其中大部分文件由于即将超时而不会进行后期处理。
  • 还有一件事。如果我在TakeWhile 之后添加后处理,它会同时运行并且完全不受限制,这也是与原始语义的变化,其中后处理是异步的,但不是并发的。但是假设它很好,我们确实想要一个并发的后处理。事实上,它将完全不受限制,如果我确实限制它(使用 Select + Merge(int)),那么对 GetFileSource 的调用将不会遵守超时,因为我们已经在第一名。再次,我必须考虑它是否对我有好处。
猜你喜欢
  • 1970-01-01
  • 2021-07-31
  • 1970-01-01
  • 2021-06-14
  • 1970-01-01
  • 2010-09-15
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多