【问题标题】:ConcurrentDictionary and ConcurrentBag for AddOrUpdate on parallel并行的 AddOrUpdate 的 ConcurrentDictionary 和 ConcurrentBag
【发布时间】:2020-12-25 11:35:36
【问题描述】:

使用 ConcurrentDictionary 和 ConcurrentBag 来添加或更新值是否正确?

基本上尝试如下,

  1. 拥有包含数百万条记录的文件并尝试处理并提取到对象。

  2. 而 entry 就像,Key-Value 对,Key=WBAN 和 Value 作为对象

     var cd = new ConcurrentDictionary<String, ConcurrentBag<Data>>();
     int count = 0;
    
     foreach (var line in File.ReadLines(path).AsParallel().WithDegreeOfParallelism(5))
     {
         var sInfo = line.Split(new char[] { ',' });
         cd.AddOrUpdate(sInfo[0], new ConcurrentBag<Data>(){ new Data()
         {
             WBAN =  sInfo[0],
                 Date = string.IsNullOrEmpty(sInfo[1]) ? "" : sInfo[1],
                 time = string.IsNullOrEmpty(sInfo[2]) ? "" : sInfo[2]
     }
         }
         ,
         (oldKey, oldValue) =>
         {
             oldValue.Add(new Data()
             {
                 WBAN = sInfo[0],
                 Date = string.IsNullOrEmpty(sInfo[1]) ? "" : sInfo[1],
                 time = string.IsNullOrEmpty(sInfo[2]) ? "" : sInfo[2]
             });
    
             return oldValue;
         }
         );
     }
    

【问题讨论】:

  • File.ReadLines(path).AsParallel().WithDegreeOfParallelism(5) - 这意味着什么?编写此并行代码毫无意义,因为您的程序是 IO 绑定的。你甚至没有使用 async-IO。你给自己制造了不必要的麻烦。
  • 没有必要并行运行,除了让你头疼之外,它几乎什么都做不了,即使你让它工作也可能运行得更慢
  • 您的代码还不必要地创建了一个新的ConcurrentBag 实例,因为ConcurrentDictionary 会运行所有工厂回调,即使会发生冲突(默认情况下它会这样做,尽管有一些解决方法)。
  • 实际上,文本文件有 2000 万条或更多记录。那么,它对并行处理有用吗?
  • 如果你花更多的时间在CPU上处理,你可能会遇到一个案例,但瓶颈会是IO,这种情况下CPU上什么也没做(相对地)。此外,由于您正在阅读行尾,并且它可能已编码,因此没有简单的方法可以并行处理(但并非不可能)。最后,我会三思而后行将 200 万条记录拉入内存,这是大型对象堆上的大量不可移动数据。听起来你想要一个数据库

标签: c# multithreading parallel.foreach concurrentdictionary


【解决方案1】:
  • 您的程序是 IO 密集型的,而不是 CPU 密集型的,因此并行处理没有任何优势。
    • 它受 IO 限制,因为您的程序必须先从文件中读取该行数据才能处理该行数据,并且一般而言计算机从存储中读取数据的速度总是比它们可以读取的慢得多处理它。
    • 由于您的程序在读取的每一行上只执行微不足道的字符串操作,我可以有把握地说,将Data 元素添加到Dictionary&lt;String,List&lt;Data&gt;&gt; 所花费的时间只是很小的一小部分。您的计算机从文本文件中读取一行所需的时间。
  • 另外,避免将File.ReadLines 用于此类程序,因为这会首先将整个文件读入内存。
    • 如果您查看我的解决方案,您会发现它使用 StreamReader 逐行读取每一行,这意味着它不需要等到它先将所有内容读入内存。

因此,要以最佳性能解析该文件,您不需要任何并发集合。

就这个:

private static readonly Char[] _sep = new Char[] { ',' }; // Declared here to ensure only a single array allocation.

public static async Task< Dictionary<String,List<Data>> > ReadFileAsync( FileInfo file )
{
    const Int32 ONE_MEGABYTE = 1 * 1024 * 1024; // Use 1MB+ sized buffers for async IO. Not smaller buffers like 1024 or 4096 as those are for synchronous IO.

    Dictionary<String,List<Data>> dict = new Dictionary<String,List<Data>>( capacity: 1024 );


    using( FileStream fs = new FileStream( path, FileAccess.Read, FileMode.Open, FileShare.Read, ONE_MEGABYTE, FileOptions.Asynchronous | FileOptions.SequentialScan ) )
    using( StreamReader rdr = new StreamReader( fs ) )
    {
        String line;
        while( ( line = await rdr.ReadLineAsync().ConfigureAwait(false) ) != null )
        {
            String[] values = line.Split( sep );
            if( values.Length < 3 ) continue;

            Data d = new Data()
            {
                WBAN = values[0],
                Date = values[1],
                time = values[2]
            };

            if( !dict.TryGetValue( d.WBAN, out List<Data> list ) )
            {
                dict[ d.WBAN ] = list = new List<Data>();
            }

            list.Add( d );
        }
    }
}

更新:假设...

假设地说,由于文件 IO(尤其是异步 FileStream IO)使用大缓冲区(在本例中为 ONE_MEGABYTE 大小的缓冲区),因此程序可以将每个缓冲区(按顺序读取)传递给并行处理器.

但问题是缓冲区内的数据不能轻易地分配给各个线程:在这种情况下因为行的长度不固定,所以仍然需要单个线程通读整个缓冲区以找出换行符的位置(技术上 可以在一定程度上并行化,这将增加大量复杂性(因为您还需要处理跨越缓冲区边界的行,或仅包含单行的缓冲区等)。

而且在这种小规模下,使用线程池和并发收集类型的开销将消除并行处理的加速,因为程序仍然很大程度上受 IO 限制。

现在,如果您有一个以千兆字节为单位的文件,其中 Data 记录的大小约为 1KB,那么我将详细演示如何做到这一点,因为在这种规模下,您可能会看到适度的性能提升。

【讨论】:

  • 有道理。并感谢您提供详细的 cmets。我可以看到,有问题的代码需要 1.09 分钟才能处理,而您在上面分享的代码是 StopWatch.Elapsed 在 00:00:08 秒内,速度非常快。
  • @Cod29 是的,当ConfigureAwait(false) 被省略时(或使用ConfigureAwait(true)),整个函数将需要更长的时间才能完成,因为线程池将等待恢复while 循环的主体一个特定的线程或上下文,这在这里是不必要的。 始终在非 UI 代码中指定 ConfigureAwait(false)stackoverflow.com/questions/13489065/…
  • File.ReadLines 没有将整个文件加载到内存中。也许您将它与File.ReadAllLines 混淆了?我也非常对使用ReadLineAsync 持怀疑态度。从my experiments 开始,此方法并不是真正的异步,只是增加了开销(导致为每一行分配一个Task&lt;string&gt;),而与同步ReadLine 相比没有任何好处。
  • @TheodorZoulias 抱歉,是的 - 我在想ReadAllLines(我忘了ReadLines 完全返回一个懒惰的IEnumerable&lt;String&gt;)。
  • @TheodorZoulias 是的,我同意 ReadLineAsync 是次优的 - 我发现如果我直接使用 FileStream.ReadAsync 然后用黑客同步 StreamReader 包装该缓冲区(并手动处理\r 位于一个缓冲区的末尾而\n 位于下一个缓冲区的开头的情况) - 但表明这对于这个问题+答案来说太过分了。
【解决方案2】:

你的想法基本上是正确的,但是在实现上存在缺陷。使用foreach 语句枚举ParallelQuery 不会导致循环内的代码并行运行。在这个阶段,并行化阶段已经完成。在您的代码中实际上没有并行工作,因为在.AsParallel().WithDegreeOfParallelism(5) 之后没有附加运算符。要并行执行循环,您必须将 foreach 替换为 ForAll 运算符,如下所示:

File.ReadLines(path)
    .AsParallel()
    .WithDegreeOfParallelism(5)
    .ForAll(line => { /* Process each line in parallel */ });

了解这里并行化的内容很重要。每一行的处理都是并行的,而从文件系统中加载每一行却不是。加载是序列化的。并行 LINQ 引擎使用的工作线程(其中一个是当前线程)在访问源 IEnumerable(在本例中为 File.ReadLines(path))时会同步。

使用嵌套的ConcurrentDictionary&lt;String, ConcurrentBag&lt;Data&gt;&gt; 结构来存储已处理的行不是很有效。您可以相信PLINQ 在分组数据方面做得比手动处理并发集合等方面做得更好。 通过使用ToLookup 运算符,您可以获得ILookup&lt;string, Data&gt;,它本质上是一个只读字典,每个键都有多个值。

var separators = new char[] { ',' };

var lookup = File.ReadLines(path)
    .AsParallel()
    .WithDegreeOfParallelism(5)
    .Select(line => line.Split(separators))
    .ToLookup(sInfo => sInfo[0], sInfo => new Data()
    {
        WBAN =  sInfo[0],
        Date = string.IsNullOrEmpty(sInfo[1]) ? "" : sInfo[1],
        time = string.IsNullOrEmpty(sInfo[2]) ? "" : sInfo[2]
    });

这在性能和内存效率方面应该是一个更好的选择,除非您出于某种原因特别希望生成的结构是可变的和线程安全的。

还有两个注意事项:

  1. 硬编码并行度(在本例中为 5)是可以的,前提是您知道运行程序的硬件。否则可能会因超额订阅而引起摩擦(线程数多于机器的实际内核数)。提示:虚拟机通常配置为单线程。

  2. ConcurrentBag 是一个非常专业的集合。在大多数情况下,使用ConcurrentQueue 可以获得更好的性能。这两个类都提供了类似的 API。人们可能更喜欢ConcurrentBag,因为它的Add 方法比Enqueue 更熟悉。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-09-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-05-16
    相关资源
    最近更新 更多