【问题标题】:Best way to convert thread safe collection to DataTable?将线程安全集合转换为 DataTable 的最佳方法?
【发布时间】:2016-03-23 20:47:52
【问题描述】:

所以这里是场景:

我必须获取一组数据,对其进行处理,构建一个对象,然后将这些对象插入到数据库中。

为了提高性能,我使用并行循环对数据进行多线程处理,并将对象存储在 CollectionBag 列表中。

那部分工作正常。但是,这里的问题是我现在需要获取该列表,将其转换为 DataTable 对象并将数据插入数据库。这很丑陋,我觉得我没有以最好的方式做到这一点(下面的伪):

ConcurrentBag<FinalObject> bag = new ConcurrentBag<FinalObject>();

ParallelOptions parallelOptions = new ParallelOptions();
parallelOptions.MaxDegreeOfParallelism = Environment.ProcessorCount;

Parallel.ForEach(allData, parallelOptions, dataObj =>
{   
    .... Process data ....

    bag.Add(theData);

    Thread.Sleep(100);
});

DataTable table = createTable();
foreach(FinalObject moveObj in bag) {
    table.Rows.Add(moveObj.x);
}

【问题讨论】:

  • 您也可以在并行循环中将 FinalObject 转换为 DataRow,以增加更多性能,将包设置为 Concurrent
  • 所以您只是将底层对象的一个​​属性添加到数据表中?如果集合中已经有了对象,为什么还需要数据表?为什么不首先填充数据表?
  • 我为这个例子简化了它 - 我使用的数据表(我使用了 9 个)范围从 7 列到 13 列
  • Nemo - 基本上有一个并发的 DataRows 包并在最后添加它们?
  • 为什么你的Parallel.ForEach 正文中会有一个Thread.Sleep

标签: c# multithreading


【解决方案1】:

这是 PLINQ(或 Rx - 我将专注于 PLINQ,因为它是基类库的一部分)的理想选择。

IEnumerable<FinalObject> bag = allData
    .AsParallel()
    .WithDegreeOfParallelism(Environment.ProcessorCount)
    .Select(dataObj =>
    {
        FinalObject theData = Process(dataObj);

        Thread.Sleep(100);

        return theData;
    });

DataTable table = createTable();

foreach (FinalObject moveObj in bag)
{
    table.Rows.Add(moveObj.x);
}

实际上,您不应通过Thread.Sleep 限制循环,而应进一步限制最大并行度,直到将 CPU 使用率降至所需水平。

免责声明:以下所有内容仅供娱乐,尽管它确实确实有效。

当然,您总是可以将它提升一个档次并生成一个完整的异步Parallel.ForEach 实现,允许您并行处理输入并异步进行节流,而不会阻塞任何线程池线程。

async Task ParallelForEachAsync<TInput, TResult>(IEnumerable<TInput> input,
                                                 int maxDegreeOfParallelism,
                                                 Func<TInput, Task<TResult>> body,
                                                 Action<TResult> onCompleted)
{
    Queue<TInput> queue = new Queue<TInput>(input);

    if (queue.Count == 0) {
        return;
    }

    List<Task<TResult>> tasksInFlight = new List<Task<TResult>>(maxDegreeOfParallelism);

    do
    {
        while (tasksInFlight.Count < maxDegreeOfParallelism && queue.Count != 0)
        {
            TInput item = queue.Dequeue();
            Task<TResult> task = body(item);

            tasksInFlight.Add(task);
        }

        Task<TResult> completedTask = await Task.WhenAny(tasksInFlight).ConfigureAwait(false);

        tasksInFlight.Remove(completedTask);

        TResult result = completedTask.GetAwaiter().GetResult(); // We know the task has completed. No need for await.

        onCompleted(result);
    }
    while (queue.Count != 0 || tasksInFlight.Count != 0);
}

用法(full Fiddle here):

async Task<DataTable> ProcessAllAsync(IEnumerable<InputObject> allData)
{
    DataTable table = CreateTable();
    int maxDegreeOfParallelism = Environment.ProcessorCount;

    await ParallelForEachAsync(
        allData,
        maxDegreeOfParallelism,
        // Loop body: these Tasks will run in parallel, up to {maxDegreeOfParallelism} at any given time.
        async dataObj =>
        {
            FinalObject o = await Task.Run(() => Process(dataObj)).ConfigureAwait(false); // Thread pool processing.

            await Task.Delay(100).ConfigureAwait(false); // Artificial throttling.

            return o;
        },
        // Completion handler: these will be executed one at a time, and can safely mutate shared state.
        moveObj => table.Rows.Add(moveObj.x)
    );

    return table;
}

struct InputObject
{
    public int x;
}

struct FinalObject
{
    public int x;
}

FinalObject Process(InputObject o)
{
    // Simulate synchronous work.
    Thread.Sleep(100);

    return new FinalObject { x = o.x };
}

相同的行为,但没有 Thread.SleepConcurrentBag&lt;T&gt;

【讨论】:

  • 感谢您的建议 - 我花时间尝试了一下,看看最终结果会给我带来什么。它快了大约 2 秒,所以它并没有做太多。但是,我正在查看您的其他一些建议,看看我能从中得到什么。谢谢!
  • @user2124871,如果您追求“更快”,那么您绝对不想要Thread.SleepTask.Delay 或任何其他类型的人为延迟。无论您使用 Parallel.ForEach 解决方案还是我的 PLINQ 版本,除了在代码样式方面,都不会真正产生太大影响。您应该分析解决方案的处理组件(“处理数据”)并对其进行优化以减少 CPU 使用率。
【解决方案2】:

听起来您通过使所有内容并行运行使事情变得相当复杂,但是如果您将DataRow obejcts 存储在您的包中而不是普通对象,最后您可以使用DataTableExtensions 创建一个DataTable 很容易来自通用集合:

var dataTable = bag.CopyToDataTable();

只需在您的项目中添加对System.Data.DataSetExtensions 的引用即可。

【讨论】:

  • CopyToDataTable&lt;T&gt; 具有T : DataRow 的通用约束,因此它似乎只适用于通用集合的一小部分(导致像这样的解决方法:msdn.microsoft.com/en-us/library/bb669096(v=vs.110).aspx)。我错过了什么吗?
  • @KirillShlenskiy 这就是为什么我说“如果您将 DataRow 对象存储在包中而不是普通对象”。我仍然对为什么 OP 使用所有这些复杂性来并行创建一个集合感到困惑,只是为了转身并将其串行转换为DataTable
  • 并行收集的原因是循环中发生的处理量。按顺序执行非常耗时,这使我可以更快地移动。数据表用于插入我需要放入数据库的所有记录。由于它不是线程安全的集合,所以我将其存储在线程安全的集合中并转换为介质以插入 DB。这可能是矫枉过正,试图弄清楚。
  • @DStanley,抱歉,我忽略了“如果您将 DataRow 对象存储在您的包中”部分。
  • 您是否尝试过将对象一一插入数据库而不是从DataTable 插入?还是您使用 SqlBulkCopy 插入记录?
【解决方案3】:

我认为这样的东西应该提供更好的性能,看起来 object[] 是比 DataRow 更好的选择,因为您需要 DataTable 来获取 DataRow 对象。

ConcurrentBag<object[]> bag = new ConcurrentBag<object[]>();

Parallel.ForEach(allData, 
    new ParallelOptions { MaxDegreeOfParallelism = Environment.ProcessorCount }, 
    dataObj =>
{
    object[] row = new object[colCount];

    //do processing

    bag.Add(row);

    Thread.Sleep(100);
});

DataTable table = createTable();
foreach (object[] row in bag)
{
    table.Rows.Add(row);
}

【讨论】:

    猜你喜欢
    • 2011-04-02
    • 1970-01-01
    • 2011-06-09
    • 1970-01-01
    • 1970-01-01
    • 2012-01-20
    • 1970-01-01
    • 1970-01-01
    • 2015-04-14
    相关资源
    最近更新 更多