【问题标题】:How to detect lines that are unique in large file using Reactive Extensions如何使用 Reactive Extensions 检测大文件中唯一的行
【发布时间】:2019-04-24 20:55:03
【问题描述】:

我必须处理大型 CSV 文件(高达数十 GB),如下所示:

Key,CompletedA,CompletedB
1,true,NULL
2,true,NULL
3,false,NULL
1,NULL,true
2,NULL,true  

我有一个解析器,可以将解析后的行生成为IEnumerable<Record>,因此我一次只能将一行读入内存。

现在我必须按 Key 对记录进行分组,并检查 CompletedA 和 CompletedB 列是否在组内具有值。在输出上,我需要记录,即组内没有 CompletedA,CompleteddB。

在这种情况下,它是用键 3 记录的。

但是,在同一个数据集上会有许多类似的处理,我不会多次迭代它。

我想我可以将 IEnumerable 转换为 IObservable 并使用 Reactive Extentions 来查找记录。

是否可以在 IObservable 集合上使用简单的 Linq 表达式以节省内存的方式进行操作?

【问题讨论】:

  • 当然,你也可以使用管道处理器,比如数据流,orrrr Reactive Extensions,但是,这都是多余的,你可以在 foreach 循环中有效地做到这一点,你会帮自己一个忙尝试这是第一个
  • records.CountBy(z => new { Key = z.Key, Value = z.CompletedA ?? z.CompletedB}).Where(z => z.Value == 1).Select(z => z.Key) 可能会帮助您入门。为此,您需要 nuget.org/packages/morelinq
  • 你有多少个不同的Keys?
  • @TheGeneral:这只是众多此类分析之一,我必须在单个 foreach 中完成所有这些分析。 foreach 不适合还有其他原因
  • @DmitryBychenko:“你有多少个不同的键?”:记录数的一半或更多。不确定生产中会有多少行,但考虑到文件大小,很多。

标签: c# system.reactive yield file-processing


【解决方案1】:

假设Key 是一个整数,我们可以尝试使用Dictionary 和一次扫描:

 // value: 0b00 - neither A nor B
 //        0b01 - A only
 //        0b10 - B only
 //        0b11 - Both A and B    
 Dictionary<int, byte> Status = new Dictionary<int, byte>();

 var query = File
   .ReadLines(@"c:\MyFile.csv")
   .Where(line => !string.IsNullOrWhiteSpace(line))
   .Skip(1) // skip header 
   .Select(line => YourParserHere(line));

 foreach (var record in query) {
   int mask = (record.CompletedA != null ? 1 : 0) |
              (record.CompletedB != null ? 2 : 0); 

   if (Status.TryGetValue(record.Key, out var value))
     Status[record.Key] = (byte) (value | mask);
   else
     Status.Add(record.Key, (byte) mask);
 }

 // All keys that don't have 3 == 0b11 value (both A and B)  
 var bothAandB = Status
   .Where(pair => pair.Value != 3)
   .Select(pair => pair.Key); 

【讨论】:

  • 我询问RX解决方案的原因是单个foreach中的东西太多,所以我需要以某种方式拆分它而不多次枚举记录,所以我认为s push collection with multiple subscribers会工作。此外,我有不同的场景和不同的“分析”。 RX 解决方案将使单个“分析”成为可重复使用的良好和平。
  • @Liero - 你需要确保IEnumerable&lt;Record&gt; 是惰性的以使 Rx 高效,但如果是这样,那么一个简单的循环也将是。
【解决方案2】:

我认为这将满足您的需求:

var result =
    source
        .GroupBy(x => x.Key)
        .SelectMany(xs =>
            (xs.Select(x => x.CompletedA).Any(x => x != null && x == true) && xs.Select(x => x.CompletedA).Any(x => x != null && x == true))
            ? new List<Record>()
            : xs.ToList());

在这里使用 Rx 没有帮助。

【讨论】:

    【解决方案3】:

    是的,Rx 库非常适合这种同步枚举一次/多次计算的操作。您可以使用Subject&lt;Record&gt; 作为一对多传播器,然后您应该将各种 Rx 运算符附加到它,然后您应该使用源可枚举的记录来提供它,最后您将从附件中收集结果现在将完成的运算符。这是基本模式:

    IEnumerable<Record> source = GetRecords();
    var subject = new Subject<Record>();
    var task1 = SomeRxTransformation1(subject);
    var task2 = SomeRxTransformation2(subject);
    var task3 = SomeRxTransformation3(subject);
    source.ToObservable().Subscribe(subject); // This line does all the work
    var result1 = task1.Result;
    var result2 = task2.Result;
    var result3 = task3.Result;
    

    SomeRxTransformation1SomeRxTransformation2 等是接受IObservable&lt;Record&gt; 并返回一些通用Task 的方法。他们的签名应该是这样的:

    Task<TResult> SomeRxTransformation1(IObservable<Record> source);
    

    例如,您想要进行的特殊分组将需要如下转换:

    Task<Record[][]> GroupByKeyExcludingSomeGroups(IObservable<Record> source)
    {
        return source
            .GroupBy(record => record.Key)
            .Select(grouped => grouped.ToArray())
            .Merge()
            .Where(array => array.All(r => !r.CompletedA && !r.CompletedB))
            .ToArray()
            .ToTask();
    }
    

    当你将它合并到模式中时,它会是这样的:

    Task<Record[][]> task1 = GroupByKeyExcludingSomeGroups(subject);
    source.ToObservable().Subscribe(subject); // This line does all the work
    Record[][] result1 = task1.Result;
    

    【讨论】:

    • 顺便说一句,用.ToObservable().Subscribe(subject) 喂给主题很简单,但效率不高。要获得性能更好的替代方案,您可以查看here
    • 另一个优化可能是实现同步ISubject,因为内置的Subject 具有嵌入的同步功能(Volatile.Read)。同步处理不需要这些,只会增加开销。实现ISubject 并不是很困难。 Here 是一个起点。
    猜你喜欢
    • 2013-11-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-03-05
    • 2014-02-16
    • 1970-01-01
    相关资源
    最近更新 更多