فصل ۲۲: Aggregate موازی، کلاس Parallel و آغاز Task Parallelism
فروش یا انتشار این ترجمه منوط به داشتن مجوز لازم از صاحب حقوق اثر است.
شکل 22-4 ـ مقایسهٔ chunk partitioning و range partitioning.
ParallelEnumerable.Range یک ParallelQuery<T> برمیگرداند؛ بنابراین پس از آن نیازی به فراخوانی AsParallel نیست.
بهینهسازی aggregationهای سفارشی
PLINQ operatorهای Sum، Average، Min و Max را بدون دخالت اضافی بهطور کارآمد موازی میکند. اما operator Aggregate برای PLINQ چالشهای ویژهای ایجاد میکند. همانطور که در فصل ۹ دیدیم، Aggregate برای aggregationهای سفارشی است. مثال زیر یک sequence عددی را جمع میکند و رفتار Sum را تقلید میکند:
int[] numbers = { 1, 2, 3 };
int sum = numbers.Aggregate (0, (total, n) => total + n); // 6
در فصل ۹ همچنین دیدیم که برای aggregation بدون seed، delegate دادهشده باید associative و commutative باشد. اگر این قانون نقض شود PLINQ نتیجهٔ نادرست میدهد، چون برای aggregate کردن چند partition بهطور همزمان، چند seed از input sequence میگیرد.
Aggregation با seed صریح در نگاه اول انتخاب امنی برای PLINQ به نظر میرسد، اما معمولاً بهصورت ترتیبی اجرا میشود، زیرا فقط یک seed مشترک دارد. برای حل این مشکل PLINQ overload دیگری از Aggregate ارائه میکند که اجازه میدهد چند seed ــ دقیقتر بگوییم یک seed factory function ــ تعریف کنید. این function برای هر thread اجرا میشود و یک seed مستقل میسازد؛ آن seed به accumulator محلی همان thread تبدیل میشود و عناصر همانجا در آن aggregate میشوند.
همچنین باید function دیگری ارائه کنید که نحوهٔ ترکیب accumulator محلی با accumulator اصلی را مشخص کند. در پایان، این overload از Aggregate یک delegate دیگر برای transform نهایی نتیجه میخواهد. چهار delegate به ترتیب عبارتاند از:
seedFactory- یک accumulator محلی جدید برمیگرداند.
updateAccumulatorFunc- یک عنصر را در accumulator محلی aggregate میکند.
combineAccumulatorFunc- accumulator محلی را با accumulator اصلی ترکیب میکند.
resultSelector- هر transform نهایی را روی نتیجه اعمال میکند.
مثال بسیار سادهٔ زیر مقادیر آرایهٔ numbers را جمع میکند:
numbers.AsParallel().Aggregate (
() => 0, // seedFactory
(localTotal, n) => localTotal + n, // updateAccumulatorFunc
(mainTot, localTot) => mainTot + localTot, // combineAccumulatorFunc
finalResult => finalResult) // resultSelector
این مثال ساختگی است، چون همان نتیجه با روشهای سادهتر مانند aggregate بدون seed یا بهتر از آن Sum به همان اندازه کارآمد به دست میآید. مثال واقعیتر: فرض کنید میخواهیم فراوانی هر حرف الفبای انگلیسی را در یک string محاسبه کنیم. نسخهٔ ترتیبی ساده چنین است:
string text = "Let’s suppose this is a really long string";
var letterFrequencies = new int[26];
foreach (char c in text)
{
int index = char.ToUpper (c) - 'A';
if (index >= 0 && index < 26) letterFrequencies [index]++;
};
برای موازیکردن این کار میتوان foreach را با Parallel.ForEach جایگزین کرد، اما در آن صورت باید concurrency روی array مشترک را خودمان مدیریت کنیم؛ و lock کردن دسترسی به array تقریباً تمام ظرفیت parallelization را از بین میبرد.
Aggregate راهحل تمیزی ارائه میکند. accumulator در اینجا همان آرایهای شبیه letterFrequencies است. نسخهٔ ترتیبی با Aggregate:
int[] result =
text.Aggregate (
new int[26], // Create the "accumulator"
(letterFrequencies, c) => // Aggregate a letter into the accumulator
{
int index = char.ToUpper (c) - 'A';
if (index >= 0 && index < 26) letterFrequencies [index]++;
return letterFrequencies;
});
و نسخهٔ موازی با overload ویژهٔ PLINQ:
int[] result =
text.AsParallel().Aggregate (
() => new int[26], // Create a new local accumulator
(localFrequencies, c) => // Aggregate into the local accumulator
{
int index = char.ToUpper (c) - 'A';
if (index >= 0 && index < 26) localFrequencies [index]++;
return localFrequencies;
},
// Aggregate local->main accumulator
(mainFreq, localFreq) =>
mainFreq.Zip (localFreq, (f1, f2) => f1 + f2).ToArray(),
finalResult => finalResult // Perform any final transformation
); // on the end result.
دقت کنید function مربوط به local accumulation آرایهٔ localFrequencies را mutate میکند. این optimization مهم و قانونی است، چون localFrequencies مخصوص همان thread است.
کلاس Parallel
PFX از طریق سه متد static در کلاس Parallel شکل پایهای از structured parallelism فراهم میکند:
Parallel.Invoke- آرایهای از delegateها را موازی اجرا میکند.
Parallel.For- معادل موازی loop نوع
for در C# را اجرا میکند. Parallel.ForEach- معادل موازی loop نوع
foreach در C# را اجرا میکند.
هر سه متد تا پایان تمام کار block میکنند. مانند PLINQ، اگر exception کنترلنشدهای رخ دهد، workerهای باقیمانده پس از iteration جاری متوقف میشوند و exception یا exceptionها در قالب AggregateException به caller برگردانده میشوند.
Parallel.Invoke
Parallel.Invoke آرایهای از delegateهای Action را موازی اجرا میکند و تا پایانشان منتظر میماند. سادهترین امضای آن:
public static void Invoke (params Action[] actions);
مانند PLINQ، متدهای Parallel.* برای کار compute-bound بهینه شدهاند نه I/O-bound. با این حال، دانلود همزمان دو web page مثال سادهای برای نمایش Parallel.Invoke است:
Parallel.Invoke (
() => new WebClient().DownloadFile ("http://www.linqpad.net", "lp.html"),
() => new WebClient().DownloadFile ("http://microsoft.com", "ms.html"));
در ظاهر این یک shortcut برای ساخت و انتظار روی دو Task وابسته به thread است؛ اما تفاوت مهمی وجود دارد: اگر آرایهای با یک میلیون delegate بدهید، Parallel.Invoke همچنان کارآمد است. دلیل آن این است که تعداد زیاد عناصر را به batchهایی تقسیم میکند که به چند Task زیرین اختصاص داده میشوند، نه اینکه برای هر delegate یک Task جدا بسازد.
مانند همهٔ متدهای Parallel، collating نتیجه بر عهدهٔ خود شماست و باید thread safety را در نظر بگیرید. کد زیر thread-unsafe است:
var data = new List<string>();
Parallel.Invoke (
() => data.Add (new WebClient().DownloadString ("http://www.foo.com")),
() => data.Add (new WebClient().DownloadString ("http://www.far.com")));
lock کردن عملیات Add مسئله را حل میکند، اما اگر تعداد زیادی delegate سریع داشته باشید، lock میتواند گلوگاه ایجاد کند. انتخاب بهتر یک collection ایمن در برابر thread است؛ در این مثال ConcurrentBag مناسب خواهد بود.
Parallel.Invoke overloadی دارد که ParallelOptions میپذیرد:
public static void Invoke (ParallelOptions options,
params Action[] actions);
با ParallelOptions میتوانید cancellation token وارد کنید، بیشینهٔ concurrency را محدود کنید و task scheduler سفارشی مشخص کنید. cancellation token زمانی مرتبط است که تعداد taskها تقریباً بیش از تعداد coreها باشد: پس از cancellation، delegateهایی که هنوز شروع نشدهاند کنار گذاشته میشوند؛ اما delegateهایی که در حال اجرا هستند تا پایان ادامه میدهند.
Parallel.For و Parallel.ForEach
این دو متد معادل loopهای for و foreach در C# را اجرا میکنند، با این تفاوت که iterationها بهجای ترتیبی، موازی اجرا میشوند. سادهترین امضاها:
public static ParallelLoopResult For (
int fromInclusive, int toExclusive, Action<int> body)
public static ParallelLoopResult ForEach<TSource> (
IEnumerable<TSource> source, Action<TSource> body)
این loop ترتیبی:
for (int i = 0; i < 100; i++)
Foo (i);
بهشکل زیر موازی میشود:
Parallel.For (0, 100, i => Foo (i));
// یا سادهتر:
Parallel.For (0, 100, Foo);
و این foreach:
foreach (char c in "Hello, world")
Foo (c);
به این شکل تبدیل میشود:
Parallel.ForEach ("Hello, world", Foo);
برای نمونه، با import فضای نام System.Security.Cryptography میتوان شش رشتهٔ keypair عمومی/خصوصی را موازی تولید کرد:
var keyPairs = new string[6];
Parallel.For (0, keyPairs.Length,
i => keyPairs[i] = RSA.Create().ToXmlString (true));
همانند Parallel.Invoke، میتوان تعداد زیادی work item به Parallel.For و Parallel.ForEach داد و آنها بهطور کارآمد روی چند Task partition میشوند. همان مثال keypair را میتوان با PLINQ نیز نوشت:
string[] keyPairs =
ParallelEnumerable.Range (0, 6)
.Select (i => RSA.Create().ToXmlString (true))
.ToArray();
loop بیرونی در برابر loop درونی
Parallel.For و Parallel.ForEach معمولاً روی loop بیرونی بهتر از loop درونی عمل میکنند، زیرا در loop بیرونی chunkهای بزرگتری از کار برای موازیسازی عرضه میشود و overhead مدیریت رقیقتر میشود. موازیکردن همزمان loop داخلی و خارجی معمولاً لازم نیست. در مثال زیر معمولاً برای سودبردن از parallelization داخلی به بیش از 100 core نیاز دارید:
Parallel.For (0, 100, i =>
{
Parallel.For (0, 50, j => Foo (i, j)); // Sequential would be better
}); // for the inner loop.
Parallel.ForEach دارای index
گاهی دانستن index هر iteration مفید است. در foreach ترتیبی ساده است، اما increment کردن یک متغیر مشترک در محیط موازی thread-safe نیست. overload زیر index را در آرگومان سوم body میدهد:
public static ParallelLoopResult ForEach<TSource> (
IEnumerable<TSource> source, Action<TSource,ParallelLoopState,long> body)
Parallel.ForEach ("Hello, world", (c, state, i) =>
{
Console.WriteLine (c.ToString() + i);
});
مثال Spellchecker با Parallel.ForEach
اگر dictionary و آرایهٔ یک میلیون واژهٔ آزمایشی مثال قبلی را بارگذاری کرده باشیم، میتوان spellcheck را با نسخهٔ indexed از Parallel.ForEach انجام داد:
var wordLookup = new HashSet<string> (
File.ReadAllLines ("WordLookup.txt"),
StringComparer.InvariantCultureIgnoreCase);
var random = new Random();
string[] wordList = wordLookup.ToArray();
string[] wordsToTest = Enumerable.Range (0, 1000000)
.Select (i => wordList [random.Next (0, wordList.Length)])
.ToArray();
wordsToTest [12345] = "woozsh";
wordsToTest [23456] = "wubsie";
var misspellings = new ConcurrentBag<Tuple<int,string>>();
Parallel.ForEach (wordsToTest, (word, state, i) =>
{
if (!wordLookup.Contains (word))
misspellings.Add (Tuple.Create ((int) i, word));
});
نتایج باید در collection ایمن در برابر thread collate شوند؛ این نقطه ضعف در مقایسه با PLINQ است. مزیت نسبت به PLINQ این است که هزینهٔ indexed Select را نداریم، چون indexed ForEach کارآمدتر است.
ParallelLoopState: خروج زودهنگام از loop
چون body در Parallel.For یا Parallel.ForEach یک delegate است، نمیتوانید با statement معمولی break از loop خارج شوید. در عوض باید روی شیء ParallelLoopState متد Break یا Stop را صدا بزنید:
public class ParallelLoopState
{
public void Break();
public void Stop();
public bool IsExceptional { get; }
public bool IsStopped { get; }
public long? LowestBreakIteration { get; }
public bool ShouldExitCurrentIteration { get; }
}
گرفتن ParallelLoopState آسان است، زیرا overloadهای For و ForEach bodyهایی از نوع Action<TSource,ParallelLoopState> میپذیرند. برای موازیکردن این کد:
foreach (char c in "Hello, world")
if (c == ',')
break;
else
Console.Write (c);
مینویسیم:
Parallel.ForEach ("Hello, world", (c, loopState) =>
{
if (c == ',')
loopState.Break();
else
Console.Write (c);
});
// OUTPUT: Hlloe
از خروجی مشخص است که bodyهای loop ممکن است با ترتیب تصادفی کامل شوند. با این تفاوت، Break دستکم همان عناصری را پوشش میدهد که اجرای ترتیبی تا نقطهٔ break پوشش میداد؛ در این مثال همیشه حروف H، e، l، l و o با ترتیبی نامشخص چاپ میشوند. در مقابل، Stop همهٔ threadها را وادار میکند بلافاصله پس از iteration جاری تمام شوند؛ بنابراین ممکن است فقط زیرمجموعهای از این حروف چاپ شود. Stop وقتی مفید است که چیزی را که دنبالش بودید پیدا کردهاید یا خطایی رخ داده و دیگر نتیجهها را بررسی نمیکنید.
اگر body طولانی باشد میتوانید در چند نقطه ShouldExitCurrentIteration را poll کنید تا در صورت Break یا Stop زودتر خارج شوید. این property بلافاصله پس از Stop و اندکی پس از Break true میشود. همچنین پس از درخواست cancellation یا رخدادن exception در loop true میشود. IsExceptional مشخص میکند آیا exception در thread دیگری رخ داده است یا نه. هر exception کنترلنشده باعث میشود loop پس از iteration فعلی هر thread متوقف شود؛ برای جلوگیری از آن باید exceptionها را صریحاً handle کنید.
بهینهسازی با local valueها
Parallel.For و Parallel.ForEach overloadهایی با type argument عمومی TLocal دارند که برای بهینهکردن collating داده در loopهای iteration-heavy طراحی شدهاند. سادهترین امضا:
public static ParallelLoopResult For <TLocal> (
int fromInclusive,
int toExclusive,
Func <TLocal> localInit,
Func <int, ParallelLoopState, TLocal, TLocal> body,
Action <TLocal> localFinally);
در عمل این overloadها بهندرت لازم میشوند، چون PLINQ بیشتر سناریوهای هدف را پوشش میدهد. مسئلهای که حل میکنند چنین است: فرض کنید میخواهیم ریشهٔ دوم اعداد 1 تا 10,000,000 را جمع کنیم. محاسبهٔ ده میلیون square root بهراحتی موازی میشود، اما جمعزدن نتیجه مشکل دارد، چون update کردن total نیازمند lock است:
object locker = new object();
double total = 0;
Parallel.For (1, 10000000,
i => { lock (locker) total += Math.Sqrt (i); });
هزینهٔ گرفتن ده میلیون lock بهعلاوهٔ blocking ایجادشده، سود parallelization را از بین میبرد.
در واقع به ده میلیون lock نیاز نداریم. تصور کنید گروهی داوطلب حجم زیادی زباله جمع میکنند. اگر همه یک سطل مشترک داشته باشند، رفتوآمد و contention فرایند را بسیار ناکارآمد میکند. راه واضح این است که هر worker یک سطل «محلی» داشته باشد و گاهی آن را در سطل اصلی خالی کند. نسخههای TLocal دقیقاً همین الگو را پیاده میکنند. برای این کار باید دو delegate اضافی بدهیم: یکی برای مقداردهی local value جدید، و دیگری برای ترکیب local aggregation با مقدار master. همچنین body بهجای void، aggregate جدید local value را برمیگرداند:
object locker = new object();
double grandTotal = 0;
Parallel.For (1, 10000000,
() => 0.0, // Initialize the local value.
(i, state, localTotal) => // Body delegate. Notice that it
localTotal + Math.Sqrt (i), // returns the new local total.
localTotal => // Add the local value
{ lock (locker) grandTotal += localTotal; } // to the master value.
);
هنوز lock لازم است، اما فقط هنگام اضافهکردن local value به grand total؛ بنابراین فرایند بهمراتب کارآمدتر میشود. همانطور که گفته شد، PLINQ اغلب برای این سناریو انتخاب خوبی است:
ParallelEnumerable.Range (1, 10000000)
.Sum (i => Math.Sqrt (i))
Task Parallelism
Task parallelism پایینترین سطح parallelization در PFX است. کلاسهای این سطح در فضای نام System.Threading.Tasks قرار دارند:
| کلاس | هدف |
|---|
Task | مدیریت یک واحد کار |
Task<TResult> | مدیریت یک واحد کار دارای مقدار بازگشتی |
TaskFactory | ساخت taskها |
TaskFactory<TResult> | ساخت taskها و continuationها با return type یکسان |
TaskScheduler | مدیریت scheduling taskها |
TaskCompletionSource | کنترل دستی workflow یک Task |
مبانی Taskها در فصل ۱۴ پوشش داده شد. در ادامهٔ فصل ۲۲، قابلیتهای پیشرفتهٔ Task برای برنامهنویسی موازی بررسی میشوند: tuning زمانبندی Task، ایجاد رابطهٔ parent/child هنگام شروع Task از داخل Task دیگر، استفادهٔ پیشرفته از continuationها و TaskFactory.