【问题标题】:Per-thread instance object in TPL Parallel.ForEachTPL Parallel.ForEach 中的每线程实例对象
【发布时间】:2017-10-12 04:14:27
【问题描述】:

是否有一种 TPL 语法允许您将对象从池中注入到任务中,这样一个对象一次只能由一个线程使用?甚至更好 - 仅由同一个线程使用?

使用示例

假设我想创建 10 个线程来打开 10 个文件:1.txt2.txt3.txt ... 10.txt 并将 500 000 个结果数字随机写入这些文件。

我可以这样做:

ConcurrentQueue<int> objs = new ConcurrentQueue<int>(); // 500000 numbers go here
Task[] tasks = Enumerable.Range(1, 10)
    .Select(i =>
    {
        return Task.Factory.StartNew(() => 
        {
            using (var f = File.Open($"{i}.txt"))
            {
                using (var wr = StreamWriter(f))
                {
                    while (objs.TryDequeue(out int obj))
                    {
                        wr.WriteLine(obj);
                    }
                }
            }
        }
    })
    .ToArray();
Task.WaitAll(tasks);

但是,是否可以在不使用并发集合的情况下仅使用 TPL 提供相同的行为?

【问题讨论】:

  • 你知道ThreadLocal吗?
  • 为什么不调整 ForEachParallelOptionsParallel.ForEach(objs, new ParallelOptions() { MaxDegreeOfParallelism = 4 }, o =&gt; {...});另一种可能是Parallel Linq:objs.AsParallel().WithDegreeOfParallelism(4).ForAll(o = &gt; ...);
  • 这听起来像XY Problem。您有一个数据库性能问题 X,并认为您将通过并行执行大量连接 (Y) 来解决它。当失败时,您会询问 Y。并行执行多个数据库查询不会解决任何性能问题。不过,这很容易使它们变得更糟。 ORM用于处理大量数据。
  • 最后,ORM既不是报告工具也不是批量导入工具。它们可以处理单个对象图,但不能处理多行。您的实际问题是什么?
  • BTW ORM 不提供抽象或设计优势。当您导入数据时,没有实体,也没有行为。 实际实体是源、行、字段、数据、转换。在将它们存储到目的地之前,您会处理一个 行,这些行会在运行中进行转换。您可以使用几 MB 的 RAM 以这种方式处理数百万行。尝试使用 ORM

标签: c# multithreading nhibernate task-parallel-library


【解决方案1】:

如果除最后两个编辑之外的所有内容都被删除会更好。

如果问题是Can you pass an object per task (not thread) when using Parallel.?答案是:是的,您可以通过接受本地状态的any of the overloads,即有一个TLocal 类型,如this one

public static ParallelLoopResult ForEach<TSource, TLocal>(
    IEnumerable<TSource> source,
    Func<TLocal> localInit,
    Func<TSource, ParallelLoopState, TLocal, TLocal> body,
    Action<TLocal> localFinally
)

Parallel.For 不使用线程。它对数据进行分区并为每个分区创建一个任务。每个任务最终都会处理一个分区的所有数据。通常,Parallel 使用的任务数量与内核数量一样多。它还使用 current 线程进行处理,这就是它看起来阻塞当前线程的原因。它没有,它开始用于处理其中一个分区。

处理本地数据的函数允许您生成初始本地值并将其传递给每个body 调用。所有带有本地数据的重载都需要body 来重新调整(可能已修改)数据,因此Parallel 本身不必存储它。这是必不可少的,因为Parallel. 可以终止和重新启动任务。如果它必须跟踪本地数据,它将无法轻松或有效地做到这一点。

对于这个特定的示例,绕过 ORM 不适合批量操作的事实,尤其是在处理数十万个对象时,localInit 应该创建一个新会话。 body 应该使用并返回该会话,而最后,localFinally 应该处理它。

var mySessionFactory
var myData=....;
Parallel.ForEach(
    myData,
    ()=>CreateSession(),
    (record,state,session)=>{
        //process the data etc.
        return session;
    },
    (session)=>session.Dispose()
);

不过还有一些警告。 NH 将更改保存在内存中,直到它们被刷新并且缓存被清除。这将产生内存问题。一种解决方案是保持计数并定期刷新数据。状态可以是 (int counter,Session session) 元组,而不是会话:

Parallel.ForEach(
    myData,
    ()=>(counter:0,session:CreateSession()),
    (record,state,localData)=>{
        var (counter,session)=localData;
        //process the data etc.
        if (counter % 1000 ==0)
        {
            session.Flush();
            session.Clear();
        }
        return (++counter,session);
    },
    data=>data.session.Dispose()
);

更好的 解决方案是提前对对象进行批处理,以便循环在 IEnumerable&lt;MyRecord[]&gt; 数组上运行,而不是 IEnumerable&lt;MyRecord&gt;。结合批处理语句,这将减少 ORM 对批量操作造成的性能损失。

编写Batch 方法并不难,但MoreLinq 已经提供了一个,可作为源代码或 NuGet 包使用:

var myBatches=myData.Batch(1000);
Parallel.ForEach(
    myBatches,
    ()=>CreateSession(),
    (records,state,session)=>{

        foreach(var record in records)
        {
            //process the data etc.
            session.Save(record);                
        }
        session.Flush();
        session.Clear();
        return session;
    },
    data=>data.session.Dispose()
);

【讨论】:

  • 感谢您的广泛而准确的回答,这正是我想要的。很抱歉用 NHibernate 示例误导您。我也更新了我的问题。
【解决方案2】:

不,没有。

最接近的解决方案是手动创建 N 个线程(使用TaskParallel.For / Parallel.ForEach)并使用ConcurrentQueue 线程安全地分发数据。

【讨论】:

  • 从 cmets 发布 Vasek 的答案作为社区答案。
  • 是的,肯定有。检查ForEach 的重载。有很多处理本地数据的重载。它们很容易被发现,因为它们使用泛型类型 TLocal 作为本地数据类型
猜你喜欢
  • 2014-12-18
  • 1970-01-01
  • 1970-01-01
  • 2020-10-11
  • 1970-01-01
  • 2021-05-27
  • 1970-01-01
  • 2012-11-03
  • 1970-01-01
相关资源
最近更新 更多