【问题标题】:Parallel for each using double threads expected in TPL每个使用 TPL 中预期的双线程并行
【发布时间】:2012-08-05 18:04:42
【问题描述】:

我将首先对我如何理解几件事情的工作进行基本解释,然后用一个 tldr 来结束这一切;如果人们只是想解决我在这里遇到的实际问题。如果我对这里的任何理解有误,请纠正我。

TPL 代表任务并行库,它是 .NET 4.0 对尝试进一步简化线程以方便开发人员使用的回答。如果您不熟悉它,(在非常基础的级别上)您启动一个新的 Task 对象并向其传递一个委托,然后该委托在从线程池获取的后台线程上运行(通过使用线程池而不是真正制作通过使用这些现有线程而不是创建和处置新线程来节省新线程、时间和资源。

据我了解,C# 中的 Parallel.ForEach 命令将为它应该执行的每个委托生成一个新线程(可能来自线程池),但可能的例外是自动执行一个或什至可能的内联如果编译器认为它们将足够快地发生以提高效率,则可以进行更多的迭代。

与我的目标最相关的背景信息:

我正在尝试制作一个启动任务的快速程序,以便与程序的其余部分同时运行。在此任务中,Parallel.ForEach 运行 3 次“迭代”。总的来说,我们希望程序现在总共运行 5 个线程(最多):1 个用于主线程,1 个用于实际任务,最多 3 个用于 Parallel.ForEach。每个线程都有自己要完成的目标(尽管 Parallel.ForEach 都有相同的目标,但其相关 itemNumber 的值不同以计算。当主线程完成所有目标时,它使用 Task.Wait() 等待在完成任务上,它也等待 Parallel.ForEach 完成。然后使用并验证这些值。

tldr;实际问题:

运行上述想法时,Parallel.ForEach 似乎初始化了两倍于我预期的 SynchronizationContexts(本质上是另一个线程的 TPL 对象)并运行所有它们,但只等待预期的数量其中。因为 Parallel.ForEach().Wait() 命令以预期的线程运行数量完成,所以 Task 也会在它认为一切都完成时完成。然后主程序发现任务已经完成,并且当它验证当前没有更多后台线程在运行时,有时剩余的 Parrallel.ForEach() 尚未完成,因此会引发错误。

已验证线程的数量与我在每次 SynchronizationContext 的 post 调用(Async 方法启动器)时打印到调试窗口中所说的相符。每个线程也由一个主线程对象引用,否则该对象计划在完成任务时被处置,但由于由于未真正预期创建的未完成线程仍然有引用,因此处置不能正确发生。

Thread testThread = Thread.CurrentThread;
Task backgroundTask = taskFactory.StartNew(() =>
{
    Thread rootTaskThread = Thread.CurrentThread;
    Assert.AreNotEqual(testThread, rootTaskThread, "First task should not inline");
    Thread.Sleep(TimeSpan.FromSeconds(2));

    Parallel.ForEach(new[] { 1, 2, 3, 4 },
       new ParallelOptions { TaskScheduler = taskFactory.Scheduler }, (int item) => {
        Thread.Sleep(TimeSpan.FromSeconds(1));
     });
});

在上面的示例中,主线程、backgroundTask 任务和 8 个 Parallel.ForEach 线程最终都存在,其中最后 9 个是在 SynchronizationContexts 上创建的。

在 SynchronizationContext 中为我的自定义重写的唯一方法是 post,如下所示:

public override void Post(SendOrPostCallback d, object state){
    Request requestOrNull = Request.ExistsForCurrentThread() ? Request.GetForCurrentThread() as Request : null;
    Request.IAsyncContextData requestData = null;

    if (requestOrNull != null){
       requestData = requestOrNull.CaptureDataForNewThreadAndIncrementReferenceCount();
    }

    Debug.WriteLine("Task started - request data " + (requestData == null ? "DOES NOT EXIST" : "EXISTS"));

    base.Post((object internalState) => {
        // Capture the spawned thread state and restore the originating thread state
        try{
            if (requestData != null){
                Request.AttachToAsynchronousContext(requestData);
            }
            d(state);
        }
        finally{
            // Restore original spawned thread state
            if (requestData != null){
            // Disposes the request if this is the last reference to it
                Request.DetachFromAsynchronousContext(requestData);
            }
        Debug.WriteLine("Task completed - request data " + (requestData == null ? "DOES NOT EXIST" : "EXISTS"));
        }
    }, state);
 }

TaskScheduler 我相信它只做它所需的基本工作:

private readonly RequestSynchronizationContext context;
private readonly ConcurrentQueue<Task> tasks = new ConcurrentQueue<Task>();

public RequestTaskScheduler(RequestSynchronizationContext synchronizationContext)
{
    this.context = synchronizationContext;
}

protected override void QueueTask(Task task){
    this.tasks.Enqueue(task);
    this.context.Post((object state) => {
        Task nextTask;
        if (this.tasks.TryDequeue(out nextTask)) 
            this.TryExecuteTask(nextTask);
    }, null);
}

protected override bool TryExecuteTaskInline(Task task, bool taskWasPreviouslyQueued){
    if (SynchronizationContext.Current == this.context)
        return this.TryExecuteTask(task);
    else
        return false;
}

protected override IEnumerable<Task> GetScheduledTasks(){
    return this.tasks.ToArray();
}

任务工厂:

public RequestTaskFactory(RequestTaskScheduler taskScheduler)
    : base(taskScheduler)
{ }

关于为什么会发生这种情况的任何想法?

【问题讨论】:

  • 您能否发布一个简短但完整的示例来展示您的行为?尤其是您的SynchronizationContext(我假设您使用的是自定义的)。因为我没有看到你描述的行为。另外,我真的不明白为什么使用Post() 超出您的预期会给您带来任何问题?您的SynchronizationContext 是否有一些依赖于此的特殊行为?
  • 您这样做是为了达到什么目的? Assert.AreNotEqual 向我建议您正在编写某种单元测试。但是,Assert.AreNotEqual 确实在验证 TPL,而不是您的代码——所以,我看不出验证第三方代码的意义。

标签: .net c#-4.0 task-parallel-library threadpool parallel.foreach


【解决方案1】:

任务本身不创建线程。如果可以的话,由 TaskScheduler 决定做什么来使操作异步。例如一些操作使用异步 IO,在这种情况下,硬件使其异步,而不是另一个工作线程。

如何调用延续取决于同步上下文。该上下文不是另一个线程,它只是抽象了可以运行操作的标准。例如,在 WPF、WinForms、Silverlight 等中,有一个 UI 同步上下文,其动作需要在 特定线程(UI 线程或主线程,以避免异常)上执行。

ForEach 将尝试 创建线程(更具体地说,它将尝试询问同步上下文以启动多个异步操作)。调度程序确实定义了它是如何做到的。如果你给它三个任务,它可能创建三个线程,或者它可能不。它决定三个并发线程是否是一件好事。例如,如果您只有两个内核,则 ForEach 不会创建两个以上的线程,因为由于上下文切换开销,这可能比使用单个线程并顺序运行代码更糟糕。

不清楚“初始化两倍的 SynchronizationContext”是什么意思。这些不是线程。你只是说它创建的线程比你预期的多吗?还是您的意思是 Post 的调用次数超出了您的预期?您的 SynchronizationContext 类基于什么? (即它的基类是什么)。基础所做的主要定义了 Post 调用的可能性。它可能觉得需要创建另一个异步操作来跟踪其他操作...您如何让调度程序使用此上下文?

SynchronizationContext 早在 TPL 之前就已经存在(首先出现在 .NET 2.0 中)。它所做的一件事是管理异步操作请求。从您的帖子中不清楚您是否理解这一点。

更新: 对 QueueTask 的第一次调用间接来自 StartNew。 对 QueueTask 的第二次调用间接来自对 ForEach 的调用 第三次调用 QueueTask 间接来自 QueueTask 中的 TryExecuteTask 接下来对 QueueTask 的 4 次调用是针对传递给 ForEach 的主体。

根据负载,QueueTask 最多可能被调用 3 次。如果我在 QueueTask 调试和中断,QueueTask 只会被调用 7 次。

此时,由于您在 QueueTask 中执行的操作与 TPL 不同(即 TryExecuteTask 正在调用额外的操作),因此很难说为什么有时会有一些额外的 QueueTask 调用。这可能来自您实现 QueueTask 的方式,因为您有效地要求调度程序从已经异步执行的任务中排队另一个异步任务。我的猜测是,这只是时机。 QueueTask 可以被如此快速地调用(因为这是通过另一个异步操作完成的),以至于 TryExecuteTask 不知道任务已经排队并强制执行任务(强制另一个调用 QueueTask)。

如果 QueueTask 实际上导致另一个对 QueueTask 的调用,因为它还没有安排它所在的任务,这就解释了为什么最多有 10 次对 QueueTask 的调用。正是 TryExeucteTask 调用导致了每个 ForEach 主体的“双重”调用......

【讨论】:

  • 我的 SynchronizationContext 直接扩展了基本的 SynchronizationContext 并且除了 post 之外没有其他方法被覆盖。是的,post 的调用次数超出了我的预期。对于 ForEach 循环中的每次迭代,都会进行 2 次 post 调用。或者至少,它的出现使得大小为 3 的 ForEach 会导致总共 7 个 post 调用,size 4 会导致 9 个 post 调用,等等......基本的 TaskScheduler 也被使用以及没有被覆盖的基本 TaskFactory方法中的任何一个。你是对的,我不知道在 TPL 之前存在 SynchronizationContext。
  • TaskScheduler 是一个抽象类;你是如何实现这些抽象方法的?另外,你是如何实现TaskFactory的?您正在做与 TPL 不同的事情。您如何设置 TaskFactory 的 Scheduler 属性?我不知道您是否使用了默认的任务调度程序(即 SynchronizationContextTaskScheduler 或 ThreadPoolTask​​Scheduler),所以看起来您是在比较苹果和橘子。
  • 添加了我对 TaskScheduler 的实现,TaskFactory 也和我添加的一样简单。至于任何 TaskFactory Scheduler 属性,那将是我不知道曾经必须做的事情?
  • @KevinMoore 我已经在我的答案中添加了细节,但我相信这是因为在QueueTask 中对TryExecuteTask 的调用太快并且迫使现有任务(已经异步运行)异步运行任务。
  • 感谢您,这对您有很大帮助!但是,当您说取决于负载时,QueueTask 最多可能会被调用 3 次——这到底会如何发生?我的意思是,最终会如何称呼它那些额外的时间(可能是间接的,然后我假设?)。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2023-03-03
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多