【问题标题】:Exception using Rx and Await to accomplish reading file line by line async使用 Rx 和 Await 完成逐行异步读取文件的异常
【发布时间】:2014-04-14 15:58:50
【问题描述】:

我正在学习使用 RX 并尝试了这个示例。但无法修复突出显示的 while 语句中发生的异常 - while(!f.EndofStream)

我想逐行读取一个巨大的文件,并且对于每一行数据,我想在不同的线程中进行一些处理(所以我使用了 ObserverOn) 我想要整个事情异步。我想使用 ReadLineAsync,因为它返回 TASK,所以我可以将它转换为 Observables 并订阅它。

我猜我首先创建的任务线程位于 Rx 线程之间。但即使我使用 currentThread 使用 Observe 和 Subscribe,我仍然无法阻止异常。想知道我是如何使用 Rx 巧妙地完成这个 Aysnc 的。

想知道整个事情是否可以做得更简单?

    static void Main(string[] args)
    {
        RxWrapper.ReadFileWithRxAsync();
        Console.WriteLine("this should be called even before the file read begins");
        Console.ReadLine();
    }

    public static async Task ReadFileWithRxAsync()
    {
        Task t = Task.Run(() => ReadFileWithRx());
        await t;
    }


    public static void ReadFileWithRx()
    {
        string file = @"C:\FileWithLongListOfNames.txt";
        using (StreamReader f = File.OpenText(file))
        {
            string line = string.Empty;
            bool continueRead = true;

            ***while (!f.EndOfStream)***
            {
                f.ReadLineAsync()
                       .ToObservable()
                       .ObserveOn(Scheduler.Default)
                       .Subscribe(t =>
                           {
                               Console.WriteLine("custom code to manipulate every line data");
                           });
            }

        }
    }

【问题讨论】:

  • 异常类型、消息和堆栈跟踪是什么?

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


【解决方案1】:

异常是InvalidOperationException - 我对 FileStream 的内部并不十分熟悉,但根据异常消息,这是因为流上有一个正在进行的异步操作而引发的。这意味着在检查EndOfStream 之前,您必须等待任何ReadLineAsync() 调用完成。

Matthew Finlay 对您的代码进行了巧妙的改造,以解决这个紧迫的问题。然而,我认为它有自己的问题——还有一个更大的问题需要研究。让我们看一下问题的基本要素:

  • 您的文件非常大。
  • 您希望异步处理它。

这表明您不希望整个文件都在内存中,您希望在处理完成时得到通知,并且可能希望尽可能快地处理文件。

两种解决方案都使用线程来处理每一行(ObserveOn 将每一行从线程池传递给一个线程)。这实际上不是一种有效的方法。

查看这两种解决方案,有两种可能性:

  • A.平均而言,读取文件行比处理文件所花费的时间更多
  • 乙。平均而言,读取文件行所花费的时间比处理它所花费的时间少。

A.文件读取一行比处理一行慢

在 A 的情况下,系统在等待文件 IO 完成时基本上会花费大部分时间处于空闲状态。在这种情况下,Matthew 的解决方案不会导致内存填满 - 但值得看看是否直接在紧密循环中使用 ReadLines 会因为线程争用较少而产生更好的结果。 (ObserveOn 将线路推到另一个线程只会在 ReadLines 没有在调用 MoveNext 之前收到线路时给你买点东西 - 我怀疑它确实 - 但测试一下!)

B.文件读取一行比处理一行快

在 B 的情况下(考虑到您的尝试,我认为更有可能),所有这些行将开始在内存中排队,并且对于足够大的文件,您最终会在内存中获得大部分.

您应该注意,除非您的处理程序触发异步代码来处理一行,否则所有行都将被串行处理,因为 Rx 保证 OnNext() 处理程序调用不会重叠。

ReadLines() 方法很棒,因为它返回一个IEnumerable<string>,并且驱动读取文件的是您对此的枚举。但是,当您对此调用 ToObservable() 时,它将尽可能快地枚举以生成可观察事件 - Rx 中没有反馈(在反应性程序中称为“背压”)来减慢此过程。

问题不在于ToObservable 本身,而在于ObserveOnObserveOn 不会阻塞 OnNext() 处理程序,它在等待订阅者处理完事件时调用它 - 它会尽可能快地将事件排队到目标调度程序。

如果您删除 ObserveOn,那么 - 只要您的 OnNext 处理程序是同步的 - 您将看到每一行都被读取和处理一次,因为 ToObservable() 正在处理相同的枚举线程作为处理程序。

如果这不是您想要的,并且您尝试通过在订阅者中触发异步作业来缓解这种情况以追求并行处理 - 例如Task.Run(() => /* process line */ 或类似的 - 那么事情不会像你希望的那样顺利。

由于处理一行比读取一行需要更长的时间,因此您将创建越来越多的任务与传入的行不同步。线程数会逐渐增加,你会饿死线程池。

在这种情况下,Rx 真的不太合适。

您可能想要的是少量工作线程(可能每个处理器内核 1 个),它们一次获取一行代码以进行处理,并限制内存中文件的行数。

一种简单的方法可能是这样,它将内存中的行数限制为固定数量的工作人员。这是一个基于拉的解决方案,在这种情况下这是一个更好的设计:

private Task ProcessFile(string filePath, int numberOfWorkers)
{
    var lines = File.ReadLines(filePath);       

    var parallelOptions = new ParallelOptions {
        MaxDegreeOfParallelism = numberOfWorkers
    };  

    return Task.Run(() => 
        Parallel.ForEach(lines, parallelOptions, ProcessFileLine));
}

private void ProcessFileLine(string line)
{
    /* Your processing logic here */
    Console.WriteLine(line);
}

并像这样使用它:

static void Main()
{       
    var processFile = ProcessFile(
        @"C:\Users\james.world\Downloads\example.txt", 8);

    Console.WriteLine("Processing file...");        
    processFile.Wait();
    Console.WriteLine("Done");
}

最后说明

有一些方法可以处理 Rx 中的背压(围绕 SO 搜索一些讨论)——但这不是 Rx 处理好的东西,而且我认为得到的解决方案比上面的替代方案可读性差。您还可以查看许多其他方法(基于角色的方法,例如 TPL 数据流,或用于高性能无锁方法的 LMAX Disruptor 样式环形缓冲区),但 pulling 的核心思想来自排队将很普遍。

即使在此分析中,我也方便地掩盖了您处理文件所做的工作,并默认每行的处理都是计算绑定的并且真正独立。如果有工作可以合并结果和/或 IO 活动来存储输出,那么所有的赌注都没有了——你也需要仔细检查这方面的效率。

在考虑并行执行工作作为优化的大多数情况下,通常有很多变量在起作用,因此最好衡量每种方法的结果以确定最佳方法。测量是一门艺术——一定要测量真实的场景,对每次测试的多次运行取平均值,并在运行之间正确重置环境(例如,消除缓存效应),以减少测量错误。

【讨论】:

  • 我认为将File.ReadAllLines替换为File.ReadLines会减少内存消耗。目前整个文件在任何事情发生之前都被读入内存。
  • 是的,你完全正确 - 这是代码示例中的一个令人讨厌的错字,老实说! :) 现在修复。谢谢。
  • @JamesWorld - 你是对的例外。感谢详细的内省。通过我的这个简单的程序读取文件,你让我思考了许多设计和性能方面的事情。干杯。
  • 不客气!顺便说一句,如果你想了解哪个线程在做什么,你可能会对my Spy method 感兴趣。另外,从线程的角度来看,this post should clarify what ObserveOn is doing
【解决方案2】:

我没有调查是什么导致了你的异常,但我认为写这个的最简洁的方法是:

File.ReadLines(file)
  .ToObservable()
  .ObserveOn(Scheduler.Default)
  .Subscribe(Console.Writeline);

注意:ReadLines 与 ReadAllLines 的不同之处在于它会在不读取整个文件的情况下开始产生,这是您想要的行为。

【讨论】:

  • 如果一行的处理时间比读入时间长,您将ToObservableObserveOn 一起使用将导致高内存消耗。删除ObserveOn 仍将导致文件的串行处理或可能的线程饥饿,具体取决于处理程序的实现。详情见我的回答。
  • @JamesWorld 欢呼,我会把它留给其他人警告
  • 这也会不必要地阻塞线程,因为ReadLines() 不是异步的。
  • @MatthewFinlay 在从ReadLines() 返回的IEnumerable 上调用MoveNext() 的那个。 (如果你不指定一个作为ToObservable 的参数,我不知道 Rx 会是哪个线程。)
  • @svick,调用MoveNext 的线程将来自TaskPool,这是我们想要完成工作的线程。
猜你喜欢
  • 2019-09-22
  • 1970-01-01
  • 1970-01-01
  • 2020-02-11
  • 2022-11-25
  • 1970-01-01
  • 2012-12-07
  • 2020-10-11
  • 1970-01-01
相关资源
最近更新 更多