【问题标题】:How to effectively perform complex processing on records read from a database如何有效地对从数据库读取的记录执行复杂的处理
【发布时间】:2019-05-14 08:29:41
【问题描述】:

公平警告:这是一个关于方法的问题,至少是良好实践...这里的问题不是语法,而是方法。

我必须非常快速地处理大量记录,并向消费者提供一组转换后的记录。我想知道是否有人对最有效的方法有实用的建议。

这是场景:

我需要执行一组相对简单的逻辑: 连接到数据库 -> 读取记录 -> 转换每条记录 -> 将输出记录提供给消费者

这个逻辑需要从一个库中获得——内部逻辑对消费者完全隐藏。 (消费者不知道发生了某种变换——他认为他只是在循环一堆对象)。

通常,我会使用这样的方法创建一个 IEnumerable 类:

public class TransformingReader<T> where T:class,new()
{
...
...
...

 public IEnumerator<T> GetEnumerator()
 {
      var items = _connection<dynamic>.GetData();
      foreach (var item in items)
      {
          T transformed = _complexTask.Transform(item);
          yield return transformed;
      }
 }
}

(这里使用动态类只是为了说明)

使用上面的类,消费者:

foreach(var item in new TransformingReader<TransactionAnalysis>())
{
    ...
    DoStuff(item);
    ...
}

事实:

  1. 我每天要处理数百万条记录 - 所以数量是个大问题。

  2. 用户 DoStuff() 函数需要一些时间才能完成。我无法预测他们的工作会有多复杂,但它肯定会比我的工作更密集。

  3. 我在一个相对受限的环境中工作 - 因此没有大量可用内存并且其他应用程序都在同一台机器上。所以,我需要负责任地行事。 (我没有在爷爷的笔记本电脑上运行——但我仍然需要编写不贪婪的合理代码)

想法:

  1. 我想尝试并行化 Transform() 函数,这样我就可以利用 DoStuff() 忙的时间来转换下一条记录。通过这种方式,希望我总是(经常?)在用户要求下一条记录时为新记录做好准备。

  2. 我希望将简单的 foreach 语法保留在消费者端。消费者无需知道我在幕后努力工作。

对于如何解决此类问题的任何想法将不胜感激。具体来说,是否有一种我不知道的模式可以帮助解决这个问题?

【问题讨论】:

    标签: c# performance optimization parallel-processing


    【解决方案1】:

    是的,这是一种生产者-消费者模式。

    请参阅Pipelines 如何实现它。

    var records = new BlockingCollection<SomeRecord>();
    var outputs = new BlockingCollection<SomeResult>();
    
    var readRecords = Task.Run(async () =>
    {
        using (var conn = new SqlConnection("..."))
        {
            conn.Open();
            using (var cmd = conn.CreateCommand())
            using (var reader = cmd.ExecuteReader())
            {
                while (reader.Read())
                {
                    var record = new SomeRecord { Prop = reader.GetValue(0) };
                    records.Add(record);
                }
            }
        }
    });
    
    var transformRecords = Task.Run(() =>
    {
        foreach (var record in records.GetConsumingEnumerable())
        {
            // transform record
            outputs.Add(new SomeResult());
        }
    });
    
    var consumeResults = Task.Run(() =>
    {
        foreach (var result in outputs.GetConsumingEnumerable())
        {
            // ...
        }
    });
    
    Task.WaitAll(readRecords, transformRecords, consumeResults);
    

    如有必要,可以轻松增加流水线阶段的数量,添加新任务。

    转换很容易并行化:

    records.GetConsumingEnumerable()
           .AsParallel()
           .AsOrdered() // if you want to keep order
    

    如果其中一个任务比其他任务快得多并且阻塞了内存,您可以限制其集合的容量:

    var records = new BlockingCollection<SomeRecord>(boundedCapacity: 50);
    

    【讨论】:

    • 这正是我所需要的,太棒了——感谢@Alexander Petrov。 MadKarel 给出的解释为更好地理解基础知识提供了背景。太可惜了,你只能选择一个答案!
    【解决方案2】:

    这听起来像Produce-consumer problem

    一种解决方案是创建一个用于检索和转换数据的线程,即生产者线程。然后在其他线程(可能是主线程)中运行consumer,用户DoStuff(item)。会有一个队列(很可能是concurrent queue)用于在线程之间进行通信。

    从用户的角度来看,您仍然可以将数据作为枚举器提供,该枚举器将从队列中读取,当队列为空时阻塞,并在它读取某个表示输入结束的预定值时结束(有时称为毒药丸)。

    内存占用由队列大小决定,因此您可以根据需要对其进行调整。

    此模式允许您扩大生产者和消费者的数量,因此您可以同时Transform() 多个项目,并同时DoStuff() 多个项目。

    根据您的描述,可以使用一个 Parallel LINQ 语句解决您的问题(在幕后使用上述解决方案的变体)。

    【讨论】:

    • 谢谢@MadKarel。太可惜了,你只能选择一个答案。这为我提供了更好地理解 Alexanders 答案的背景。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-01-16
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多