【问题标题】:TParallel.For performanceTParallel.For 性能
【发布时间】:2014-12-17 21:20:00
【问题描述】:

给定以下在一维数组中查找奇数的简单任务:

begin
  odds := 0;
  Ticks := TThread.GetTickCount;
  for i := 0 to MaxArr-1 do
      if ArrXY[i] mod 2 = 0 then
        Inc(odds);
  Ticks := TThread.GetTickCount - Ticks;
  writeln('Serial: ' + Ticks.ToString + 'ms, odds: ' + odds.ToString);
end;

看起来这将是并行处理的理想选择。因此,人们可能会想使用以下 TParallel.For 版本:

begin
  odds := 0;
  Ticks := TThread.GetTickCount;
  TParallel.For(0,  MaxArr-1, procedure(I:Integer)
  begin
    if ArrXY[i] mod 2 = 0 then
      inc(odds);
  end);
  Ticks := TThread.GetTickCount - Ticks;
  writeln('Parallel - false odds: ' + Ticks.ToString + 'ms, odds: ' + odds.ToString);
end;

这种并行计算的结果在两个方面有点令人惊讶:

  1. 计算的赔率数有误

  2. 执行时间比串口版本长

1) 是可以解释的,因为我们没有保护并发访问的几率变量。所以为了解决这个问题,我们应该改用TInterlocked.Increment(odds);。

2) 也是可以解释的:表现出false sharing的效果。

理想情况下,错误共享问题的解决方案是使用局部变量来存储中间结果,并且仅在所有并行任务结束时汇总这些中间结果。 这是我无法理解的真正问题:有没有办法将局部变量放入我的匿名方法中?请注意,在匿名方法体中简单地声明一个局部变量是行不通的,因为每次迭代都会调用匿名方法体。如果这在某种程度上可行,是否有办法在每次任务迭代结束时从匿名方法中获取我的中间结果?

编辑:我实际上对计算赔率或埃文并不感兴趣。我只是用这个来演示效果。

为了完整起见,这里有一个控制台应用程序来展示效果:

program Project4;

{$APPTYPE CONSOLE}

{$R *.res}

uses
  System.SysUtils, System.Threading, System.Classes, System.SyncObjs;

const
  MaxArr = 100000000;

var
  Ticks: Cardinal;
  i: Integer;
  odds: Integer;
  ArrXY: array of Integer;

procedure FillArray;
var
  i: Integer;
  j: Integer;
begin
  SetLength(ArrXY, MaxArr);
  for i := 0 to MaxArr-1 do
      ArrXY[i]:=Random(MaxInt);
end;

procedure Parallel;
begin
  odds := 0;
  Ticks := TThread.GetTickCount;
  TParallel.For(0,  MaxArr-1, procedure(I:Integer)
  begin
    if ArrXY[i] mod 2 = 0 then
      TInterlocked.Increment(odds);
  end);
  Ticks := TThread.GetTickCount - Ticks;
  writeln('Parallel: ' + Ticks.ToString + 'ms, odds: ' + odds.ToString);
end;

procedure ParallelFalseResult;
begin
  odds := 0;
  Ticks := TThread.GetTickCount;
  TParallel.For(0,  MaxArr-1, procedure(I:Integer)
  begin
    if ArrXY[i] mod 2 = 0 then
      inc(odds);
  end);
  Ticks := TThread.GetTickCount - Ticks;
  writeln('Parallel - false odds: ' + Ticks.ToString + 'ms, odds: ' + odds.ToString);
end;

procedure Serial;
begin
  odds := 0;
  Ticks := TThread.GetTickCount;
  for i := 0 to MaxArr-1 do
      if ArrXY[i] mod 2 = 0 then
        Inc(odds);
  Ticks := TThread.GetTickCount - Ticks;
  writeln('Serial: ' + Ticks.ToString + 'ms, odds: ' + odds.ToString);
end;

begin
  try
    FillArray;
    Serial;
    ParallelFalseResult;
    Parallel;
  except
    on E: Exception do
      Writeln(E.ClassName, ': ', E.Message);
  end;
  Readln;
end.

【问题讨论】:

  • 设置和调用所有这些匿名方法比执行方法花费的时间要多得多。所以“虚假分享”并不是这里真正的问题。正如您所发现的,必须使用 interlockedincrement,这也会使进程停止。至于存储中间体,您可以使用全局数组。不过,在这种情况下,更倾向于使用普通的单线程解决方案。
  • Map reduce 是你想要的。尽管这项任务如此琐碎,但线程开销将占主导地位
  • @LURD,如何从我的匿名方法中访问全局数组?你能举个例子吗?
  • 你有索引,只需用它来存储每次迭代的结果。我假设您的实际任务与此示例不同。或遵循 Davids 的建议,Is there a MapReduce library for Delphi?。
  • @LURD 事实上,我根本不是真正的任务。您可以将我的问题视为学术问题。抱歉,我如何使用索引来存储我的结果?你能告诉我一些代码吗,因为我不明白。

标签: multithreading delphi parallel-processing delphi-xe7


【解决方案1】:

这个问题的关键是正确的分区和尽可能少的共享。

使用此代码,它的运行速度几乎是串行代码的 4 倍。

const 
  WorkerCount = 4;

function GetWorker(index: Integer; const oddsArr: TArray<Integer>): TProc;
var
  min, max: Integer;
begin
  min := MaxArr div WorkerCount * index;
  if index + 1 < WorkerCount then
    max := MaxArr div WorkerCount * (index + 1) - 1
  else
    max := MaxArr - 1;
  Result :=
    procedure
    var
      i: Integer;
      odds: Integer;
    begin
      odds := 0;
      for i := min to max do
        if Odd(ArrXY[i]) then
          Inc(odds);
      oddsArr[index] := odds;
    end;
end;

procedure Parallel;
var
  i: Integer;
  oddsArr: TArray<Integer>;
  workers: TArray<ITask>;
begin
  odds := 0;
  Ticks := TThread.GetTickCount;
  SetLength(oddsArr, WorkerCount);
  SetLength(workers, WorkerCount);

  for i := 0 to WorkerCount-1 do
    workers[i] := TTask.Run(GetWorker(i, oddsArr));
  TTask.WaitForAll(workers);

  for i := 0 to WorkerCount-1 do
    Inc(odds, oddsArr[i]);
  Ticks := TThread.GetTickCount - Ticks;
  writeln('Parallel: ' + Ticks.ToString + 'ms, odds: ' + odds.ToString);
end;

您可以使用 TParallel.For 编写类似的代码,但它的运行速度仍然比仅使用 TTask 慢一些(例如比串行快 3 倍)。

顺便说一句,我使用该函数返回工作人员 TProc 以获得正确的索引捕获。如果您在同一例程中的循环中运行它,您将捕获循环变量。

19.12.2014 更新:

由于我们发现关键是正确的分区,因此可以很容易地将其放入并行 for 循环中,而无需将其锁定在特定的数据结构上:

procedure ParallelFor(lowInclusive, highInclusive: Integer;
  const iteratorRangeEvent: TProc<Integer, Integer>);

  procedure CalcPartBounds(low, high, count, index: Integer;
    out min, max: Integer);
  var
    len: Integer;
  begin
    len := high - low + 1;
    min := (len div count) * index;
    if index + 1 < count then
      max := len div count * (index + 1) - 1
    else
      max := len - 1;
  end;

  function GetWorker(const iteratorRangeEvent: TProc<Integer, Integer>;
    min, max: Integer): ITask;
  begin
    Result := TTask.Run(
      procedure
      begin
        iteratorRangeEvent(min, max);
      end)
  end;

var
  workerCount: Integer;
  workers: TArray<ITask>;
  i, min, max: Integer;
begin
  workerCount := TThread.ProcessorCount;
  SetLength(workers, workerCount);
  for i := 0 to workerCount - 1 do
  begin
    CalcPartBounds(lowInclusive, highInclusive, workerCount, i, min, max);
    workers[i] := GetWorker(iteratorRangeEvent, min, max);
  end;
  TTask.WaitForAll(workers);
end;

procedure Parallel4;
begin
  odds := 0;
  Ticks := TThread.GetTickCount;
  ParallelFor(0, MaxArr-1,
    procedure(min, max: Integer)
    var
      i, n: Integer;
    begin
      n := 0;
      for i := min to max do
        if Odd(ArrXY[i]) then
          Inc(n);
      AtomicIncrement(odds, n);
    end);
  Ticks := TThread.GetTickCount - Ticks;
  writeln('ParallelEx: Stefan Glienke ' + Ticks.ToString + ' ms, odds: ' + odds.ToString);
end;

关键是使用局部变量进行计数,最后只使用共享变量一次来添加小计。

【讨论】:

  • 非常好。事实上,我使用 Tasks 有一个类似的解决方案,虽然不是那么优雅。但我看不出如何将这样的解决方案硬塞到 TParallel.For 解决方案中 - 你能详细说明一下吗?
  • 将一个简单的for-to循环放入parallel.for本身会带来开销,因为每次迭代都是一个匿名方法调用(以及更多),这会破坏您示例中简单事物的性能。但要完成这项工作,我猜你需要像他们在 .NET TPL 中所说的分区器之类的东西。 TParallel.For 并不是使所有事物都平行的灵丹妙药。如本例所示,还有其他(更好的)方法可以做到这一点。
  • “正如本例所示,还有其他(更好的)方法可以做到这一点”......当我问这个问题时,这几乎是我的感受。我只是觉得我忽略了一些东西。我完全同意,分区器+聚合器将是要走的路。好吧,也许 XE8 或 9
  • 谢谢斯特凡!受您的回答启发,我想出了一个更可重用的解决方案 - 请参阅下面的答案。也许你想看看它,我相信它可以在很多方面得到改进。
  • @iamjoosy 太棒了!但是,我忍不住要对其进行一些改进,以免与数据结构耦合。查看我的编辑。
【解决方案2】:

使用来自 SVN 的 OmniThreadLibrary(尚未包含在任何正式版本中),您可以以不需要对共享计数器进行互锁访问的方式编写它。

function CountParallelOTL: integer;
var
  counters: array of integer;
  numCores: integer;
  i: integer;
begin
  numCores := Environment.Process.Affinity.Count;
  SetLength(counters, numCores);
  FillChar(counters[0], Length(counters) * SizeOf(counters[0]), 0);

  Parallel.For(0, MaxArr - 1)
    .NumTasks(numCores)
    .Execute(
      procedure(taskIndex, value: integer)
      begin
        if Odd(ArrXY[value]) then
          Inc(counters[taskIndex]);
      end);

  Result := counters[0];
  for i := 1 to numCores - 1 do
    Inc(Result, counters[i]);
end;

然而,这仍然充其量与顺序循环相当,最坏的情况是慢几倍。

我已将此与 Stefan 的解决方案(XE7 任务)和一个简单的 XE7 Parallel.For 与互锁增量(XE7 for)进行了比较。

我的具有 4 个超线程内核的笔记本的结果:

序列:543 毫秒内发现 49999640 个奇数元素

并行 (OTL):在 555 毫秒内发现 49999640 个奇数元素

并行(XE7 任务):在 136 毫秒内发现 49999640 个奇数元素

并行(XE7 for):在 1667 毫秒内发现 49999640 个奇数元素

我的带有 12 个超线程内核的工作站的结果:

序列:685 毫秒内发现 50005291 个奇数元素

并行 (OTL):在 1309 毫秒内发现 50005291 个奇数元素

并行(XE7 任务):在 62 毫秒内找到 50005291 个奇数元素

并行(XE7 for):在 3379 毫秒内找到 50005291 个奇数元素

与 System.Threading Paralell.For 相比有了很大的改进,因为没有互锁增量,但手工制作的解决方案要快得多。

完整的测试程序:

program ParallelCount;

{$APPTYPE CONSOLE}

{$R *.res}

uses
  System.SyncObjs,
  System.Classes,
  System.SysUtils,
  System.Threading,
  DSiWin32,
  OtlCommon,
  OtlParallel;

const
  MaxArr = 100000000;

var
  Ticks: Cardinal;
  i: Integer;
  odds: Integer;
  ArrXY: array of Integer;

procedure FillArray;
var
  i: Integer;
  j: Integer;
begin
  SetLength(ArrXY, MaxArr);
  for i := 0 to MaxArr-1 do
    ArrXY[i]:=Random(MaxInt);
end;

function CountSerial: integer;
var
  odds: integer;
begin
  odds := 0;
  for i := 0 to MaxArr-1 do
      if Odd(ArrXY[i]) then
        Inc(odds);
  Result := odds;
end;

function CountParallelOTL: integer;
var
  counters: array of integer;
  numCores: integer;
  i: integer;
begin
  numCores := Environment.Process.Affinity.Count;
  SetLength(counters, numCores);
  FillChar(counters[0], Length(counters) * SizeOf(counters[0]), 0);

  Parallel.For(0, MaxArr - 1)
    .NumTasks(numCores)
    .Execute(
      procedure(taskIndex, value: integer)
      begin
        if Odd(ArrXY[value]) then
          Inc(counters[taskIndex]);
      end);

  Result := counters[0];
  for i := 1 to numCores - 1 do
    Inc(Result, counters[i]);
end;

function GetWorker(index: Integer; const oddsArr: TArray<Integer>; workerCount: integer): TProc;
var
  min, max: Integer;
begin
  min := MaxArr div workerCount * index;
  if index + 1 < workerCount then
    max := MaxArr div workerCount * (index + 1) - 1
  else
    max := MaxArr - 1;
  Result :=
    procedure
    var
      i: Integer;
      odds: Integer;
    begin
      odds := 0;
      for i := min to max do
        if Odd(ArrXY[i]) then
          Inc(odds);
      oddsArr[index] := odds;
    end;
end;

function CountParallelXE7Tasks: integer;
var
  i: Integer;
  oddsArr: TArray<Integer>;
  workers: TArray<ITask>;
  workerCount: integer;
begin
  workerCount := Environment.Process.Affinity.Count;
  odds := 0;
  Ticks := TThread.GetTickCount;
  SetLength(oddsArr, workerCount);
  SetLength(workers, workerCount);

  for i := 0 to workerCount-1 do
    workers[i] := TTask.Run(GetWorker(i, oddsArr, workerCount));
  TTask.WaitForAll(workers);

  for i := 0 to workerCount-1 do
    Inc(odds, oddsArr[i]);
  Result := odds;
end;

function CountParallelXE7For: integer;
var
  odds: integer;
begin
  odds := 0;
  TParallel.For(0,  MaxArr-1, procedure(I:Integer)
  begin
    if Odd(ArrXY[i]) then
      TInterlocked.Increment(odds);
  end);
  Result := odds;
end;

procedure Count(const name: string; func: TFunc<integer>);
var
  time: int64;
  cnt: integer;
begin
  time := DSiTimeGetTime64;
  cnt := func();
  time := DSiElapsedTime64(time);
  Writeln(name, ': ', cnt, ' odd elements found in ', time, ' ms');
end;

begin
  try
    FillArray;

    Count('Serial', CountSerial);
    Count('Parallel (OTL)', CountParallelOTL);
    Count('Parallel (XE7 tasks)', CountParallelXE7Tasks);
    Count('Parallel (XE7 for)', CountParallelXE7For);

    Readln;
  except
    on E: Exception do
      Writeln(E.ClassName, ': ', E.Message);
  end;
end.

【讨论】:

  • 在这个特定示例中真正影响性能的是为每个项目调用匿名方法。我认为这是我的解决方案胜过你的唯一原因。
  • 是的,这很可能是主要原因。简单或性能 - 不能兼得...
【解决方案3】:

我想我们之前讨论过关于 OmniThreadLibrary 的问题。多线程解决方案时间较长的主要原因是TParallel.For的开销与实际计算所需的时间相比。

局部变量在这里没有任何帮助,而全局 threadvar 可能会解决错误共享问题。唉,你可能找不到在完成循环后总结所有这些踏步的方法。

IIRC,最好的方法是将任务分成合理的部分,并为每次迭代处理一系列数组条目,并增加一个专用于该部分的变量。仅此一项并不能解决错误共享问题,因为即使有不同的变量,如果它们恰好是同一缓存行的一部分,也会发生这种情况。

另一种解决方案是编写一个以串行方式处理给定数组切片的类,并行处理此类的多个实例,然后评估结果。

顺便说一句:您的代码不计算赔率 - 它计算偶数。

还有:有一个名为Odd 的内置函数,通常比您使用的mod 代码性能更好。

【讨论】:

  • 确实,Uwe,我们之前一直在一起调查这个问题。虽然您是对的,但在我的简单示例中,调用开销是主要问题,即使您增加循环中的计算负载,结果也是相似的。虽然我们可以在 Omnithread 库中成功解决该问题,但我看不到使用 TParallel.for 循环执行此操作的方法,因为它目前已实现
  • 判断一个整数是否为奇数的最快方法是读取最右边位的状态。如果为1则为奇数,为0则为偶数。我相信名为 Odd 的函数使用了这种方法,但我不确定。
  • 顺便说一句,局部变量可以解决问题!实际上,我使用带有局部变量的四个 ITask(我的机器中的 4 个核心)编写了相同的任务,并且它几乎完美地扩展了。
  • 一组 ITask 是与 TParallel.For 完全不同的方法,因此允许其他实现。问题发生在尝试将 TParallel.For 用于多线程方法只是因为单线程包含一个 for 循环。在转向多线程时,更改算法的结构是很常见的。
  • 你是绝对正确的 Uwe:TParallel.For 不是我特定任务的工具,尽管只是看串行版本可能会这么想。
【解决方案4】:

好的,受 Stefan Glienke 的回答启发,我起草了一个更可重用的 TParalleEx 类,而不是 ITasks 使用 IFutures。该类在某种程度上也模仿了带有聚合委托的 C# TPL。这只是初稿,但展示了如何相对轻松地扩展现有的 PPL。这个版本现在可以在我的系统上完美扩展 - 如果其他人可以在不同的配置上测试它,我会很高兴。感谢大家富有成效的回答和 cmets。

program Project4;

{$APPTYPE CONSOLE}

{$R *.res}

uses
  System.SysUtils, System.Threading, System.Classes, System.SyncObjs;

const
  MaxArr = 100000000;

var
  Ticks: Cardinal;
  i: Integer;
  odds: Integer;
  ArrXY: TArray<Integer>;

type

TParallelEx<TSource, TResult> = class
  private
    class function GetWorker(body: TFunc<TArray<TSource>, Integer, Integer, TResult>; source: TArray<TSource>; min, max: Integer): TFunc<TResult>;
  public
    class procedure &For(source: TArray<TSource>;
                         body: TFunc<TArray<TSource>, Integer, Integer, TResult>;
                         aggregator: TProc<TResult>);
  end;

procedure FillArray;
var
  i: Integer;
  j: Integer;
begin
  SetLength(ArrXY, MaxArr);
  for i := 0 to MaxArr-1 do
      ArrXY[i]:=Random(MaxInt);
end;

procedure Parallel;
begin
  odds := 0;
  Ticks := TThread.GetTickCount;
  TParallel.For(0,  MaxArr-1, procedure(I:Integer)
  begin
    if ArrXY[i] mod 2 <> 0 then
      TInterlocked.Increment(odds);
  end);
  Ticks := TThread.GetTickCount - Ticks;
  writeln('Parallel: ' + Ticks.ToString + 'ms, odds: ' + odds.ToString);
end;

procedure Serial;
begin
  odds := 0;
  Ticks := TThread.GetTickCount;
  for i := 0 to MaxArr-1 do
      if ArrXY[i] mod 2 <> 0 then
        Inc(odds);
  Ticks := TThread.GetTickCount - Ticks;
  writeln('Serial: ' + Ticks.ToString + 'ms, odds: ' + odds.ToString);
end;

const
  WorkerCount = 4;

function GetWorker(index: Integer; const oddsArr: TArray<Integer>): TProc;
var
  min, max: Integer;
begin
  min := MaxArr div WorkerCount * index;
  if index + 1 < WorkerCount then
    max := MaxArr div WorkerCount * (index + 1) - 1
  else
    max := MaxArr - 1;
  Result :=
    procedure
    var
      i: Integer;
      odds: Integer;
    begin
      odds := 0;
      for i := min to max do
        if ArrXY[i] mod 2 <> 0 then
          Inc(odds);
      oddsArr[index] := odds;
    end;
end;

procedure Parallel2;
var
  i: Integer;
  oddsArr: TArray<Integer>;
  workers: TArray<ITask>;
begin
  odds := 0;
  Ticks := TThread.GetTickCount;
  SetLength(oddsArr, WorkerCount);
  SetLength(workers, WorkerCount);

  for i := 0 to WorkerCount-1 do
    workers[i] := TTask.Run(GetWorker(i, oddsArr));
  TTask.WaitForAll(workers);

  for i := 0 to WorkerCount-1 do
    Inc(odds, oddsArr[i]);
  Ticks := TThread.GetTickCount - Ticks;
  writeln('Parallel: Stefan Glienke ' + Ticks.ToString + ' ms, odds: ' + odds.ToString);
end;

procedure parallel3;
var
  sum: Integer;
begin
  Ticks := TThread.GetTickCount;
  TParallelEx<Integer, Integer>.For( ArrXY,
     function(Arr: TArray<Integer>; min, max: Integer): Integer
      var
        i: Integer;
        res: Integer;
      begin
        res := 0;
        for i := min to max do
          if Arr[i] mod 2 <> 0 then
            Inc(res);
        Result := res;
      end,
      procedure(res: Integer) begin sum := sum + res; end );
  Ticks := TThread.GetTickCount - Ticks;
  writeln('ParallelEx: Markus Joos ' + Ticks.ToString + ' ms, odds: ' + odds.ToString);
end;

{ TParallelEx<TSource, TResult> }

class function TParallelEx<TSource, TResult>.GetWorker(body: TFunc<TArray<TSource>, Integer, Integer, TResult>; source: TArray<TSource>; min, max: Integer): TFunc<TResult>;
begin
  Result := function: TResult
  begin
    Result := body(source, min, max);
  end;
end;

class procedure TParallelEx<TSource, TResult>.&For(source: TArray<TSource>;
  body: TFunc<TArray<TSource>, Integer, Integer, TResult>;
  aggregator: TProc<TResult>);
var
  I: Integer;
  workers: TArray<IFuture<TResult>>;
  workerCount: Integer;
  min, max: integer;
  MaxIndex: Integer;
begin
  workerCount := TThread.ProcessorCount;
  SetLength(workers, workerCount);
  MaxIndex := length(source);
  for I := 0 to workerCount -1 do
  begin
    min := (MaxIndex div WorkerCount) * I;
    if I + 1 < WorkerCount then
      max := MaxIndex div WorkerCount * (I + 1) - 1
    else
      max := MaxIndex - 1;
    workers[i]:= TTask.Future<TResult>(GetWorker(body, source, min, max));
  end;
  for i:= 0 to workerCount-1 do
  begin
    aggregator(workers[i].Value);
  end;
end;

begin
  try
    FillArray;
    Serial;
    Parallel;
    Parallel2;
    Parallel3;
  except
    on E: Exception do
      Writeln(E.ClassName, ': ', E.Message);
  end;
  Readln;
end.

【讨论】:

    【解决方案5】:

    关于使用局部变量收集总和然后在最后收集它们的任务,您可以为此使用单独的数组:

    var
      sums: array of Integer;
    begin
      SetLength(sums, MaxArr);
      for I := 0 to MaxArr-1 do
        sums[I] := 0;
    
      Ticks := TThread.GetTickCount;
      TParallel.For(0, MaxArr-1,
        procedure(I:Integer)
        begin
          if ArrXY[i] mod 2 = 0 then
            Inc(sums[I]);
        end
      );
      Ticks := TThread.GetTickCount - Ticks;
    
      odds := 0;
      for I := 0 to MaxArr-1 do
        Inc(odds, sums[i]);
    
      writeln('Parallel - false odds: ' + Ticks.ToString + 'ms, odds: ' + odds.ToString);
    end;
    

    【讨论】:

    • 最终的串行求和遍历一个与输入数组大小相同的数组!
    • @DavidHeffernan,这里的例子不是真的。 OP 希望存储每次迭代的单独结果。
    • 我只是在演示如何做到这一点。我没有说应该这样做。
    • 首先,我认为您的解决方案很好,因为由于 Tparalell.for 中内置的非常聪明的自动调整步幅机制,sum 数组可以避免虚假共享效果。另一方面,遍历整个 sum 数组来求和结果既不是一个很好的解决方案,也不是一个非常高效的解决方案。
    • @LURD 不,他要数,那是一个reduce步骤。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-01-06
    • 1970-01-01
    • 1970-01-01
    • 2011-03-11
    • 1970-01-01
    • 2019-09-01
    • 2018-12-03
    相关资源
    最近更新 更多