【问题标题】:Parallel.Foreach with localFinally gets stalled despite completing all iterations尽管完成了所有迭代,但带有 localFinally 的 Parallel.Foreach 仍然停滞不前
【发布时间】:2012-06-07 10:50:21
【问题描述】:

在我的 Parallel.ForEach 循环中,localFinally 委托确实在所有线程上被调用。 我发现当我的并行循环停止时会发生这种情况。 在我的并行循环中,我有大约三个条件检查阶段,它们在循环完成之前返回。似乎是当线程从这些阶段返回而不是整个主体的执行时,它才不会执行 localFinally 委托。

循环结构如下:

 var startingThread = Thread.CurrentThread;
 Parallel.ForEach(fullList, opt,
         ()=> new MultipleValues(),
         (item, loopState, index, loop) =>
         {
            if (cond 1)
                return loop;
            if (cond 2)
                {
                process(item);
                return loop;
                }
            if (cond 3)
                return loop;

            Do Work(item);
            return loop;
          },
          partial =>
           {
              Log State of startingThread and threads
            } );

我在一个小数据集上运行循环并详细记录,发现虽然 Parallel.ForEach 完成了所有迭代,并且 localFinally 的最后一个线程的日志是 - 线程 6 Loop Indx 16 的调用线程状态是 WaitSleepJoin
循环仍然没有优雅地完成并且仍然停滞不前......任何线索为什么会停滞?

干杯!

【问题讨论】:

  • 某处可能出现死锁
  • 可能只是伪代码,但在当前状态下永远不会达到 cond 3。 if(cond2) 在条件周围没有括号(因此只有 process(item) 属于它)。
  • @RobertVerpalen 不,这只是由于省略了括号而导致的伪代码中的错误......
  • 快一点,日志机制是如何实现的?日志记录机制可能是问题吗?我曾经花了一天时间寻找一个问题,我的稍微错误的日志告诉我存在(但没有)=P
  • 也许您只是误解了localFinally 的含义。它不是为每个项目调用的,而是为Paralle.ForEach() 使用的每个线程调用的。并且许多项目可以共享同一个线程。

标签: c# multithreading task-parallel-library


【解决方案1】:

在看到 localFinally 的定义(在每个线程完成后执行)后进行了快速测试运行,这让我怀疑这可能意味着并行创建的线程比执行的循环少得多。例如

        var test = new List<List<string>> ();
        for (int i = 0; i < 1000; i++)
        {
            test.Add(null);
        }

        int finalcount = 0;
        int itemcount = 0;
        int loopcount = 0;

        Parallel.ForEach(test, () => new List<string>(),
            (item, loopState, index, loop) =>
            {
                Interlocked.Increment(ref loopcount);
                loop.Add("a");
                //Thread.Sleep(100);
                return loop;
            },
            l =>
            {
                Interlocked.Add(ref itemcount, l.Count);                    
                Interlocked.Increment(ref finalcount);                    
            });

在这个循环结束时,itemcount 和 loopcount 是预期的 1000,并且(在我的机器上)finalcount 为 1 或 2,具体取决于执行速度。在有条件的情况下:直接返回时,执行速度可能要快得多,并且不需要额外的线程。只有在执行任务时才需要更多线程。但是参数(在我的例子中是 l)包含所有执行的组合列表。 这可能是导致日志记录差异的原因吗?

【讨论】:

  • 您不应该使用Interlocked.IncrementInterlocked.Add 方法来避免各种计数器可能出现的线程竞争情况吗?
  • 在我的本地测试场景中使它们易变,因此它们应该是线程安全的并且结果是相同的,但是对于发布的示例,您是对的,应该使用某种锁定
  • @RobertVerpalen 将字段标记为volatile 并不会神奇地使其成为线程安全的。特别是,这并不意味着++ 可以正常工作,因为++ 不是原子的。
  • 嗯,然后似乎 volatile 还不够,只是快速检查线程睡眠以查看线程数增加,我的循环数为 987。更改代码以使用互锁功能。感谢您的提醒
  • @svick,所以我注意到,你和克里斯是绝对正确的!从现在开始,我会牢记这一点。为了清楚起见:结果保持不变。
【解决方案2】:

我想你只是误解了localFinally 的含义。它不是为每个项目调用的,而是为Parallel.ForEach() 使用的每个线程调用的。并且许多项目可以共享同一个线程。

它存在的原因是你可以在每个线程上独立地执行一些聚合,并且只在最后将它们连接在一起。这样,您只需在非常小的一段代码中处理同步(并让它影响您的性能)。

例如,如果你想计算一组项目的总和,你可以这样做:

int totalSum = 0;
Parallel.ForEach(
    collection, item => Interlocked.Add(ref totalSum, ComputeScore(item)));

但是在这里,您为每个项目都调用Interlocked.Add(),这可能会很慢。使用localInitlocalFinally,您可以像这样重写代码:

int totalSum = 0;
Parallel.ForEach(
    collection,
    () => 0,
    (item, state, localSum) => localSum + ComputeScore(item),
    localSum => Interlocked.Add(ref totalSum, localSum));

请注意,代码仅在localFinally 中使用Interlocked.Add(),并在body 中访问全局状态。这样,同步成本只需支付几次,每个线程使用一次。

注意:我在这个例子中使用了Interlocked,因为它非常简单而且很明显是正确的。如果代码比较复杂,我会先用lock,只有在需要良好性能时才尝试使用Interlocked

【讨论】:

  • 非常感谢您的启发性响应...在我的实现中,除了将它用于调试目的之外,localFinally 没有任何用处,因为我的并行。ForEach 循环执行所有迭代,但仍然停滞不前。这是 localFinally 的日志输出,它在调用循环之前打印 CurrentThread 的状态——调用线程状态是 WaitSleepJoin for Thread 6 Loop Indx 16,但尽管完成了所有迭代,但 Parallel.ForEach 没有正常退出任何线索为什么循环不会正常终止?
  • 非常感谢您的回复!我设法找到代码在其中一个调用中阻塞的点,我让循环继续完成!干杯!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-12-18
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-12-06
  • 1970-01-01
相关资源
最近更新 更多