فصل ۲۲: Task Parallelism، AggregateException و Concurrent Collections
فروش یا انتشار این ترجمه منوط به داشتن مجوز لازم از صاحب حقوق اثر است.
ساخت و شروع Taskها
همانطور که در فصل ۱۴ توضیح داده شد، Task.Run یک Task یا Task<TResult> میسازد و شروع میکند. این متد در واقع shortcut برای Task.Factory.StartNew است که بهکمک overloadهای بیشتر انعطاف بالاتری میدهد.
مشخصکردن state object
Task.Factory.StartNew اجازه میدهد state objectی مشخص کنید که به target پاس داده شود. در این حالت signature متد target باید یک parameter از نوع object داشته باشد:
var task = Task.Factory.StartNew (Greet, "Hello");
task.Wait(); // Wait for task to complete.
void Greet (object state) { Console.Write (state); } // Hello
این روش هزینهٔ closure موردنیاز برای lambdaای که Greet را صدا میزند حذف میکند. این micro-optimization در عمل بهندرت لازم است؛ بنابراین استفادهٔ مفیدتر state object، اختصاص یک نام معنادار به Task است. سپس میتوان نام را از AsyncState خواند:
var task = Task.Factory.StartNew (state => Greet ("Hello"), "Greeting");
Console.WriteLine (task.AsyncState); // Greeting
task.Wait();
void Greet (string message) { Console.Write (message); }
TaskCreationOptions
هنگام فراخوانی StartNew یا instantiate کردن Task میتوانید با enum پرچمی TaskCreationOptions نحوهٔ اجرا را تنظیم کنید. مقدارهای قابل ترکیب مهم عبارتاند از:
LongRunning, PreferFairness, AttachedToParent
LongRunning به scheduler پیشنهاد میدهد یک thread اختصاصی برای Task در نظر بگیرد. همانطور که در فصل ۱۴ گفته شد، این گزینه برای Taskهای I/O-bound و کارهای طولانی مفید است تا Taskهای کوتاه مجبور نشوند زمان نامعقولی در صف بمانند.
PreferFairness از scheduler میخواهد تا حد امکان Taskها را مطابق ترتیب شروعشدن زمانبندی کند. scheduler در حالت عادی ممکن است برای optimization از local work-stealing queueها استفاده کند؛ این روش ساخت child taskها را بدون contention ناشی از یک work queue مشترک ممکن میکند.
Child Taskها
وقتی یک Task از داخل Task دیگری شروع میشود، میتوان رابطهٔ parent-child ایجاد کرد. child task با TaskCreationOptions.AttachedToParent ساخته میشود:
Task parent = Task.Factory.StartNew (() =>
{
Console.WriteLine ("I am a parent");
Task.Factory.StartNew (() => // Detached task
{
Console.WriteLine ("I am detached");
});
Task.Factory.StartNew (() => // Child task
{
Console.WriteLine ("I am a child");
}, TaskCreationOptions.AttachedToParent);
});
ویژگی child task این است که انتظار برای completion parent، انتظار برای تمام childها را نیز شامل میشود و exceptionهای child در همان نقطه به بالا منتقل میشوند:
TaskCreationOptions atp = TaskCreationOptions.AttachedToParent;
var parent = Task.Factory.StartNew (() =>
{
Task.Factory.StartNew (() => // Child
{
Task.Factory.StartNew (() => { throw null; }, atp); // Grandchild
}, atp);
});
// The following call throws a NullReferenceException (wrapped
// in nested AggregateExceptions):
parent.Wait();
انتظار برای چند Task
برای یک Task میتوان از Wait یا در Task<TResult> از Result استفاده کرد. برای چند Task، متدهای static Task.WaitAll و Task.WaitAny وجود دارند. WaitAll تا پایان همهٔ Taskها منتظر میماند و نسبت به انتظار جداگانه برای هر Task efficientتر است، چون حداکثر به یک context switch نیاز دارد. اگر چند Task exception بدهند، WaitAll همچنان برای همه صبر میکند و سپس AggregateExceptionای حاوی exceptionهای Taskهای faulted پرتاب میکند؛ در عمل معادل الگویی شبیه این است:
// Assume t1, t2 and t3 are tasks:
var exceptions = new List<Exception>();
try { t1.Wait(); } catch (AggregateException ex) { exceptions.Add (ex); }
try { t2.Wait(); } catch (AggregateException ex) { exceptions.Add (ex); }
try { t3.Wait(); } catch (AggregateException ex) { exceptions.Add (ex); }
if (exceptions.Count > 0) throw new AggregateException (exceptions);
WaitAny از نظر مفهوم شبیه انتظار روی ManualResetEventSlimای است که هر Task هنگام completion آن را signal میکند. علاوه بر timeout، میتوان cancellation token به متدهای Wait داد تا خود عملیات انتظار cancel شود، نه Task.
Cancel کردن Taskها
میتوانید هنگام شروع Task یک cancellation token بدهید. اگر cancellation از طریق همان token رخ دهد، Task به state «Canceled» میرود:
var cts = new CancellationTokenSource();
CancellationToken token = cts.Token;
cts.CancelAfter (500);
Task task = Task.Factory.StartNew (() =>
{
Thread.Sleep (1000);
token.ThrowIfCancellationRequested(); // Check for cancellation request
}, token);
try { task.Wait(); }
catch (AggregateException ex)
{
Console.WriteLine (ex.InnerException is TaskCanceledException); // True
Console.WriteLine (task.IsCanceled); // True
Console.WriteLine (task.Status); // Canceled
}
TaskCanceledException subclassی از OperationCanceledException است. اگر بخواهید خودتان OperationCanceledException پرتاب کنید، باید cancellation token را به constructor آن بدهید. در غیر این صورت Task به TaskStatus.Canceled نمیرود و continuationهای OnlyOnCanceled فعال نمیشوند.
اگر Task پیش از شروع cancel شود اصلاً schedule نمیشود و بلافاصله OperationCanceledException روی Task ایجاد میشود. چون cancellation tokenها توسط APIهای دیگر نیز شناخته میشوند، میتوان آنها را به سازههای دیگر پاس داد تا cancellation بهصورت پیوسته propagate شود:
var cancelSource = new CancellationTokenSource();
CancellationToken token = cancelSource.Token;
Task task = Task.Factory.StartNew (() =>
{
// Pass our cancellation token into a PLINQ query:
var query = someSequence.AsParallel().WithCancellation (token)...
... enumerate query ...
});
فراخوانی Cancel روی cancelSource، PLINQ query را cancel میکند؛ query یک OperationCanceledException در body Task پرتاب میکند و همین Task را نیز cancel میکند.
Continuationها
متد ContinueWith بلافاصله پس از پایان یک Task، delegateای را اجرا میکند:
Task task1 = Task.Factory.StartNew (() => Console.Write ("antecedent.."));
Task task2 = task1.ContinueWith (ant => Console.Write ("..continuation"));
بهمحض اینکه task1 ــ antecedent ــ کامل شود، fail شود یا cancel شود، task2 ــ continuation ــ شروع میشود. اگر antecedent قبل از رسیدن اجرای برنامه به خط دوم تمام شده باشد، continuation بلافاصله schedule میشود. آرگومان ant reference به Task antecedent است. خود ContinueWith نیز Task برمیگرداند و chain کردن continuationهای بیشتر را ساده میکند.
بهطور پیشفرض antecedent و continuation ممکن است روی threadهای متفاوت اجرا شوند. با TaskContinuationOptions.ExecuteSynchronously میتوان آنها را روی همان thread اجرا کرد؛ این کار برای continuationهای بسیار ریزدانه با کاهش indirection میتواند performance را بهتر کند.
Continuation و Task<TResult>
continuation نیز میتواند از نوع Task<TResult> باشد و داده برگرداند. مثال زیر مقدار Math.Sqrt(8*2) را با زنجیرهای از Taskها محاسبه میکند:
Task.Factory.StartNew<int> (() => 8)
.ContinueWith (ant => ant.Result * 2)
.ContinueWith (ant => Math.Sqrt (ant.Result))
.ContinueWith (ant => Console.WriteLine (ant.Result)); // 4
Continuationها و exceptionها
continuation میتواند با property به نام Exception در antecedent بفهمد Task قبلی fault شده است؛ یا کافی است Result/Wait را فراخوانی کند و AggregateException حاصل را بگیرد. اگر antecedent fault شود و continuation هیچیک از این کارها را نکند، exception «unobserved» تلقی میشود و هنگامی که Task بعداً garbage-collected شود event static به نام TaskScheduler.UnobservedTaskException رخ میدهد.
الگوی امن، دوباره پرتابکردن exceptionهای antecedent است. مادامی که روی continuation انتظار انجام شود، exception propagate شده و به waiter برگردانده میشود:
Task continuation = Task.Factory.StartNew (() => { throw null; })
.ContinueWith (ant =>
{
ant.Wait();
// Continue processing...
});
continuation.Wait(); // Exception is now thrown back to caller.
راه دیگر تعیین continuationهای جداگانه برای نتیجهٔ exceptional و nonexceptional با TaskContinuationOptions است:
Task task1 = Task.Factory.StartNew (() => { throw null; });
Task error = task1.ContinueWith (ant => Console.Write (ant.Exception),
TaskContinuationOptions.OnlyOnFaulted);
Task ok = task1.ContinueWith (ant => Console.Write ("Success!"),
TaskContinuationOptions.NotOnFaulted);
این الگو بهخصوص همراه child taskها مفید است. extension method زیر exceptionهای کنترلنشدهٔ Task را «میبلعد»:
public static void IgnoreExceptions (this Task task)
{
task.ContinueWith (t => { var ignore = t.Exception; },
TaskContinuationOptions.OnlyOnFaulted);
}
Task.Factory.StartNew (() => { throw null; }).IgnoreExceptions();
میتوان آن را با logging بهبود داد.
Continuationها و child taskها
ویژگی قدرتمند continuation این است که فقط وقتی اجرا میشود که همهٔ child taskها پایان یافته باشند. در آن نقطه exceptionهای child نیز به continuation marshal میشوند. مثال زیر سه child task را شروع میکند که هرکدام NullReferenceException میدهند و سپس همه را با یک continuation روی parent دریافت میکند:
TaskCreationOptions atp = TaskCreationOptions.AttachedToParent;
Task.Factory.StartNew (() =>
{
Task.Factory.StartNew (() => { throw null; }, atp);
Task.Factory.StartNew (() => { throw null; }, atp);
Task.Factory.StartNew (() => { throw null; }, atp);
})
.ContinueWith (p => Console.WriteLine (p.Exception),
TaskContinuationOptions.OnlyOnFaulted);
شکل 22-5 ـ continuation پس از پایان parent و تمام childهای attached اجرا میشود و exceptionهای آنها را دریافت میکند.
Continuationهای شرطی
بهطور پیشفرض continuation بدون شرط schedule میشود؛ چه antecedent موفق تمام شود، exception بدهد یا cancel شود. با پرچمهای قابل ترکیب TaskContinuationOptions میتوان رفتار را تغییر داد. سه flag اصلی:
NotOnRanToCompletion = 0x10000,
NotOnFaulted = 0x20000,
NotOnCanceled = 0x40000,
این flagها subtractive هستند؛ هرچه بیشتر اعمال شوند احتمال اجرای continuation کمتر میشود. برای راحتی، مقدارهای ترکیبی زیر نیز وجود دارند:
OnlyOnRanToCompletion = NotOnFaulted | NotOnCanceled,
OnlyOnFaulted = NotOnRanToCompletion | NotOnCanceled,
OnlyOnCanceled = NotOnRanToCompletion | NotOnFaulted
ترکیب هر سه Not* بیمعناست، چون continuation همیشه cancel خواهد شد. RanToCompletion یعنی antecedent بدون cancellation یا exception کنترلنشده موفق بوده است. Faulted یعنی antecedent exception کنترلنشده داده است. Canceled یا یعنی antecedent با cancellation token خودش cancel شده و OperationCanceledException متناظر رخ داده، یا antecedent بهطور ضمنی به این دلیل cancel شده که predicate یک conditional continuation را برآورده نکرده است.
نکتهٔ کلیدی این است که اگر continuation بهسبب این flagها اجرا نشود، فراموش یا abandon نمیشود؛ بلکه state آن Canceled میشود. بنابراین continuationهای روی همان continuation، مگر با NotOnCanceled محدود شوند، همچنان اجرا میشوند:
Task t1 = Task.Factory.StartNew (...);
Task fault = t1.ContinueWith (ant => Console.WriteLine ("fault"),
TaskContinuationOptions.OnlyOnFaulted);
Task t3 = fault.ContinueWith (ant => Console.WriteLine ("t3"));
در این حالت t3 همیشه schedule میشود؛ حتی اگر t1 exception ندهد. اگر t1 موفق باشد، Task به نام fault cancel میشود و چون روی t3 محدودیتی نیست، t3 بیشرط اجرا میشود. برای آنکه فقط در صورت اجرای واقعی fault، t3 اجرا شود:
Task t3 = fault.ContinueWith (ant => Console.WriteLine ("t3"),
TaskContinuationOptions.NotOnCanceled);
میتوان بهجای آن OnlyOnRanToCompletion داد؛ تفاوت این است که در آن حالت اگر خود fault exception بدهد، t3 اجرا نمیشود.
شکل 22-6 ـ رفتار conditional continuationها و تأثیر canceled شدن continuation میانی.
Continuation با چند antecedent
با ContinueWhenAll و ContinueWhenAny در TaskFactory میتوان continuation را بر پایهٔ completion چند antecedent schedule کرد. پس از معرفی task combinatorهای WhenAll و WhenAny در فصل ۱۴، این متدها تا حد زیادی زائد شدهاند. مثال:
var task1 = Task.Run (() => Console.Write ("X"));
var task2 = Task.Run (() => Console.Write ("Y"));
var continuation = Task.Factory.ContinueWhenAll (
new[] { task1, task2 }, tasks => Console.WriteLine ("Done"));
// Same result with WhenAll:
var continuation2 = Task.WhenAll (task1, task2)
.ContinueWith (ant => Console.WriteLine ("Done"));
چند continuation روی یک antecedent
فراخوانی چندبارهٔ ContinueWith روی یک Task، چند continuation روی یک antecedent میسازد. پس از پایان antecedent همهٔ continuationها با هم شروع میشوند، مگر اینکه ExecuteSynchronously مشخص شده باشد که در آن صورت sequential اجرا میشوند. مثال زیر یک ثانیه صبر میکند و سپس XY یا YX مینویسد:
var t = Task.Factory.StartNew (() => Thread.Sleep (1000));
t.ContinueWith (ant => Console.Write ("X"));
t.ContinueWith (ant => Console.Write ("Y"));
Task Schedulerها
Task scheduler، Taskها را به threadها تخصیص میدهد و با کلاس abstract TaskScheduler نمایش داده میشود. .NET دو implementation concrete فراهم میکند: scheduler پیشفرض که با CLR thread pool کار میکند، و synchronization context scheduler. دومی عمدتاً برای مدل threading در WPF و Windows Forms طراحی شده است؛ جایی که UI elementها و controlها باید فقط از thread سازندهشان دسترسی داده شوند.
با capture کردن scheduler مربوط به UI میتوان Task یا continuation را وادار کرد روی همان context اجرا شود:
// Suppose we are on a UI thread in a Windows Forms / WPF application:
_uiScheduler = TaskScheduler.FromCurrentSynchronizationContext();
Task.Run (() => Foo())
.ContinueWith (ant => lblResult.Content = ant.Result, _uiScheduler);
البته برای چنین کاری معمولاً async functionهای C# انتخاب رایجتری هستند. نوشتن task scheduler سفارشی با subclass کردن TaskScheduler ممکن است، اما فقط در سناریوهای بسیار تخصصی؛ برای custom scheduling معمولاً TaskCompletionSource مناسبتر است.
TaskFactory
Task.Factory property staticای روی Task است که یک TaskFactory پیشفرض برمیگرداند. هدف factory ساخت سه نوع Task است: Taskهای معمولی با StartNew، continuationهای دارای چند antecedent با ContinueWhenAll/ContinueWhenAny، و Taskهایی که متدهای الگوی قدیمی APM را با FromAsync میپیچند.
ساخت TaskFactory سفارشی
TaskFactory یک abstract factory نیست و میتوان آن را instantiate کرد. این کار زمانی مفید است که بارها Taskهایی با مقادیر غیرپیشفرض یکسان برای TaskCreationOptions، TaskContinuationOptions یا TaskScheduler میسازید:
var factory = new TaskFactory (
TaskCreationOptions.LongRunning | TaskCreationOptions.AttachedToParent,
TaskContinuationOptions.None);
Task task1 = factory.StartNew (Method1);
Task task2 = factory.StartNew (Method2);
...
گزینههای continuation سفارشی هنگام ContinueWhenAll و ContinueWhenAny اعمال میشوند.
کار با AggregateException
PLINQ، کلاس Parallel و Taskها exceptionها را بهصورت خودکار به consumer marshal میکنند. ضرورت این رفتار را میتوان با LINQ query زیر دید که در iteration نخست DivideByZeroException میدهد:
try
{
var query = from i in Enumerable.Range (0, 1000000)
select 100 / i;
...
}
catch (DivideByZeroException)
{
...
}
اگر PLINQ این query را موازی کند و handling exception را نادیده بگیرد، DivideByZeroException احتمالاً روی thread دیگری رخ میدهد، catch ما را دور میزند و application را از کار میاندازد. بنابراین exceptionها بهصورت خودکار گرفته و برای caller دوباره پرتاب میشوند.
اما چون چند thread بهکار میرود، ممکن است همزمان دو یا چند exception رخ دهد. برای گزارش همه، exceptionها در containerی از نوع AggregateException بستهبندی میشوند. property به نام InnerExceptions همهٔ exceptionهای گرفتهشده را نگه میدارد:
try
{
var query = from i in ParallelEnumerable.Range (0, 1000000)
select 100 / i;
// Enumerate query
...
}
catch (AggregateException aex)
{
foreach (Exception ex in aex.InnerExceptions)
Console.WriteLine (ex.Message);
}
Flatten و Handle
Flatten
AggregateExceptionها اغلب خود شامل AggregateExceptionهای دیگری هستند؛ مثلاً وقتی child task exception میدهد. برای حذف تمام سطحهای nesting، Flatten را صدا بزنید. این متد AggregateException جدیدی با فهرست تخت inner exceptionها برمیگرداند:
catch (AggregateException aex)
{
foreach (Exception ex in aex.Flatten().InnerExceptions)
myLogWriter.LogException (ex);
}
Handle
گاهی میخواهید فقط typeهای مشخصی از exception را handle کنید و بقیه دوباره پرتاب شوند. متد Handle یک predicate از نوع زیر میگیرد و آن را روی همهٔ inner exceptionها اجرا میکند:
public void Handle (Func<Exception, bool> predicate)
اگر predicate برای exceptionی true برگرداند، آن exception handled محسوب میشود. پس از بررسی همهٔ exceptionها، اگر همه handled باشند exception دیگری پرتاب نمیشود؛ اگر برخی false باشند، AggregateException جدیدی فقط با exceptionهای unhandled ساخته و پرتاب میشود.
مثال زیر در نهایت AggregateException دیگری شامل فقط یک NullReferenceException پرتاب میکند:
var parent = Task.Factory.StartNew (() =>
{
// We’ll throw 3 exceptions at once using 3 child tasks:
int[] numbers = { 0 };
var childFactory = new TaskFactory
(TaskCreationOptions.AttachedToParent, TaskContinuationOptions.None);
childFactory.StartNew (() => 5 / numbers[0]); // Division by zero
childFactory.StartNew (() => numbers [1]); // Index out of range
childFactory.StartNew (() => { throw null; }); // Null reference
});
try { parent.Wait(); }
catch (AggregateException aex)
{
aex.Flatten().Handle (ex =>
{
if (ex is DivideByZeroException)
{
Console.WriteLine ("Divide by zero");
return true;
}
if (ex is IndexOutOfRangeException)
{
Console.WriteLine ("Index out of range");
return true;
}
return false;
});
}
Concurrent Collections
.NET در فضای نام System.Collections.Concurrent collectionهای thread-safe زیر را ارائه میکند:
| Concurrent collection | معادل غیرهمزمان |
|---|
ConcurrentStack<T> | Stack<T> |
ConcurrentQueue<T> | Queue<T> |
ConcurrentBag<T> | ندارد |
ConcurrentDictionary<TKey,TValue> | Dictionary<TKey,TValue> |
Concurrent collectionها برای سناریوهای concurrency بالا بهینه شدهاند، اما هر جا collection thread-safe لازم باشد میتوانند جایگزین lock دور collection معمولی شوند. چند نکته:
- Collectionهای معمولی در همهٔ سناریوها جز concurrency بسیار بالا سریعترند.
- Thread-safe بودن collection تضمین نمیکند کدی که از آن استفاده میکند thread-safe باشد.
- اگر هنگام تغییر collection توسط thread دیگر آن را enumerate کنید، exception رخ نمیدهد؛ در عوض ترکیبی از محتوای قدیمی و جدید میبینید.
- نسخهٔ concurrent برای
List<T> وجود ندارد. - Concurrent stack، queue و bag درون خود از linked list استفاده میکنند. این ساختار از
Stack/Queue معمولی حافظهٔ بیشتری مصرف میکند، اما برای دسترسی concurrent مناسبتر است، چون linked list پیادهسازی lock-free یا low-lock را ساده میکند؛ افزودن node فقط چند reference را تغییر میدهد، در حالی که insertion در ساختار شبیه List<T> ممکن است هزاران عنصر را جابهجا کند.
بنابراین این collectionها صرفاً shortcut برای lock کردن collection معمولی نیستند. برای نمونه، کد زیر روی یک thread:
var d = new ConcurrentDictionary<int,int>();
for (int i = 0; i < 1000000; i++) d[i] = 123;
حدود سه برابر کندتر از این است:
var d = new Dictionary<int,int>();
for (int i = 0; i < 1000000; i++) lock (d) d[i] = 123;
با این حال خواندن از ConcurrentDictionary سریع است، چون readها lock-free هستند.
Concurrent collectionها همچنین متدهای خاصی برای عملیات atomic test-and-act مانند TryPop دارند. بیشتر این متدها زیر interface به نام IProducerConsumerCollection<T> یکپارچه شدهاند.
IProducerConsumerCollection<T>
Producer/consumer collection مجموعهای است که دو use case اصلی دارد: افزودن عنصر یا «producing»، و گرفتن عنصر همراه با حذف آن یا «consuming». مثال کلاسیک stack و queue است. این collectionها در parallel programming مهماند، چون با implementationهای lock-free کارآمد سازگارند.
interface IProducerConsumerCollection<T> یک producer/consumer collection thread-safe را نمایش میدهد و کلاسهای ConcurrentStack<T>، ConcurrentQueue<T> و ConcurrentBag<T> آن را پیادهسازی میکنند. این interface از ICollection گسترش یافته و متدهای زیر را اضافه میکند:
void CopyTo (T[] array, int index);
T[] ToArray();
bool TryAdd (T item);
bool TryTake (out T item);
TryAdd و TryTake ابتدا بررسی میکنند عملیات add/remove قابل انجام است یا نه و در صورت امکان همان عملیات را انجام میدهند. test و act بهصورت atomic انجام میشوند و نیازی به lock معمولی نیست:
int result;
lock (myStack) if (myStack.Count > 0) result = myStack.Pop();
TryTake اگر collection خالی باشد false برمیگرداند. TryAdd در سه implementation استاندارد همیشه موفق است؛ اما اگر concurrent collection سفارشیای بنویسید که duplicate را ممنوع کند، میتوانید در صورت وجود element false برگردانید.
عنصری که TryTake حذف میکند به subclass بستگی دارد: در stack تازهترین عنصر، در queue قدیمیترین عنصر، و در bag هر عنصری که بتوان efficientتر حذف کرد. کلاسهای concrete عمدتاً TryTake/TryAdd را explicitly implement میکنند و همان قابلیت را با نامهای public اختصاصی مثل TryDequeue و TryPop ارائه میدهند.
ConcurrentBag<T>
ConcurrentBag<T> مجموعهای بدون ترتیب از objectها نگه میدارد و duplicate را مجاز میداند. زمانی مناسب است که هنگام Take/TryTake مهم نیست کدام element دریافت شود.
مزیت bag نسبت به concurrent queue یا stack این است که متد Add هنگام فراخوانی همزمان توسط threadهای زیاد تقریباً contention ندارد؛ در queue یا stack هنوز مقداری contention وجود دارد، هرچند بسیار کمتر از lock دور collection معمولی. Take در bag نیز بسیار efficient است، مادامی که هر thread بیش از تعداد elementهایی که خودش Add کرده Take نکند.
درون concurrent bag، هر thread linked list خصوصی خود را دارد. element به list thread فراخوانندهٔ Add اضافه میشود و contention حذف میشود. هنگام enumeration، enumerator از private listهای threadها عبور میکند. هنگام Take ابتدا list همان thread بررسی میشود؛ اگر element داشته باشد، عملیات بدون contention انجام میشود. اگر خالی باشد، bag باید از list thread دیگری element «بدزدد» و اینجا contention ممکن است رخ دهد.
بنابراین بهطور دقیق، Take تازهترین element اضافهشده روی همان thread را میدهد و اگر چنین عنصری نباشد، تازهترین element یکی از threadهای دیگر را که بهطور تصادفی انتخاب میشود برمیگرداند.
Concurrent bag زمانی ایدهآل است که عملیات موازی بیشتر شامل Add باشد یا Add و Take روی هر thread متعادل باشند. برای producer/consumer queue که producer و consumer threadهای متفاوتاند انتخاب خوبی نیست.
BlockingCollection<T>
اگر TryTake روی ConcurrentStack<T>، ConcurrentQueue<T> یا ConcurrentBag<T> و collection خالی باشد، false برمیگرداند. گاهی بهتر است در این حالت تا رسیدن element منتظر بمانیم. طراحان PFX بهجای افزودن overloadهای متعدد timeout/cancellation، این قابلیت را در wrapperی به نام BlockingCollection<T> قرار دادند.
Blocking collection هر IProducerConsumerCollection<T> را wrap میکند و اجازه میدهد element را با Take بردارید؛ اگر element موجود نباشد call block میشود. همچنین میتوان اندازهٔ کل collection را محدود کرد؛ اگر bound پر شود producer نیز block میشود. چنین مجموعهای bounded blocking collection نام دارد.
روش استفاده:
- کلاس را instantiate کنید و در صورت نیاز collection زیرین و بیشینهٔ اندازه را مشخص کنید.
- با
Add یا TryAdd element اضافه کنید. - با
Take یا TryTake element مصرف کنید.
اگر constructor بدون collection فراخوانی شود، بهطور خودکار ConcurrentQueue<T> ساخته میشود. متدهای producer/consumer از cancellation token و timeout پشتیبانی میکنند. Add/TryAdd در collection bounded ممکن است block شوند؛ Take/TryTake هنگام خالیبودن block میشوند.
روش دیگر مصرف، GetConsumingEnumerable است که sequence بالقوه بینهایتی برمیگرداند و elementها را هنگام آمادهشدن yield میکند. با CompleteAdding میتوان sequence را خاتمه داد و افزودن element جدید را نیز ممنوع کرد.
BlockingCollection متدهای static AddToAny و TakeFromAny دارد که با دادن چند blocking collection، add یا take را روی نخستین collectionی که قادر به پاسخگویی است انجام میدهند.
نوشتن Producer/Consumer Queue
Producer/consumer queue هم در برنامهنویسی موازی و هم در concurrency عمومی ساختار مفیدی است:
- queue شامل work itemها یا دادهای است که باید روی آن کار انجام شود.
- وقتی کاری باید اجرا شود enqueue میشود و caller به کارهای دیگر ادامه میدهد.
- یک یا چند worker thread در background itemها را از queue برمیدارند و اجرا میکنند.
این الگو کنترل دقیقی روی تعداد worker threadهای فعال میدهد و نه فقط CPU، بلکه resourceهای دیگر را نیز محدود میکند. مثلاً برای disk I/O سنگین میتوان concurrency را محدود کرد تا OS و applicationهای دیگر starve نشوند. workerها را میتوان در طول عمر queue بهصورت dynamic اضافه یا کم کرد. خود CLR thread pool نوعی producer/consumer queue است که برای jobهای کوتاه compute-bound بهینه شده است.
معمولاً queue شامل دادههایی است که همان task روی آنها انجام میشود، مثلاً filenameهایی که باید encrypt شوند. اگر item را delegate در نظر بگیریم، queue عمومیتری میسازیم که هر item میتواند کار متفاوتی انجام دهد.
نمونهٔ زیر با BlockingCollection<Action> پیادهسازی سادهای ارائه میکند:
public class PCQueue : IDisposable
{
BlockingCollection<Action> _taskQ = new BlockingCollection<Action>();
public PCQueue (int workerCount)
{
// Create and start a separate Task for each consumer:
for (int i = 0; i < workerCount; i++)
Task.Factory.StartNew (Consume);
}
public void Enqueue (Action action) { _taskQ.Add (action); }
void Consume()
{
// This sequence that we’re enumerating will block when no elements
// are available and will end when CompleteAdding is called.
foreach (Action action in _taskQ.GetConsumingEnumerable())
action(); // Perform task.
}
public void Dispose() { _taskQ.CompleteAdding(); }
}
چون چیزی به constructor BlockingCollection ندادیم، یک concurrent queue بهطور خودکار ساخته شد. اگر ConcurrentStack میدادیم، producer/consumer stack داشتیم.
استفاده از Taskها در Queue
نسخهٔ قبلی انعطاف کمی دارد، چون پس از enqueue نمیتوان work item را track کرد. بهتر است بتوانیم completion را بفهمیم و await کنیم، work item را cancel کنیم و exceptionهای آن را تمیز مدیریت کنیم. شیء مناسب همین Task است؛ میتوان آن را با TaskCompletionSource یا مستقیم بهصورت Task شروعنشده ساخت:
public class PCQueue : IDisposable
{
BlockingCollection<Task> _taskQ = new BlockingCollection<Task>();
public PCQueue (int workerCount)
{
for (int i = 0; i < workerCount; i++)
Task.Factory.StartNew (Consume);
}
public Task Enqueue (Action action, CancellationToken cancelToken
= default (CancellationToken))
{
var task = new Task (action, cancelToken);
_taskQ.Add (task);
return task;
}
public Task<TResult> Enqueue<TResult> (Func<TResult> func,
CancellationToken cancelToken = default (CancellationToken))
{
var task = new Task<TResult> (func, cancelToken);
_taskQ.Add (task);
return task;
}
void Consume()
{
foreach (var task in _taskQ.GetConsumingEnumerable())
try
{
if (!task.IsCanceled) task.RunSynchronously();
}
catch (InvalidOperationException) { } // Race condition
}
public void Dispose() { _taskQ.CompleteAdding(); }
}
در Enqueue Task را میسازیم اما شروع نمیکنیم، آن را queue میکنیم و به caller برمیگردانیم. در Consume Task بهصورت synchronous روی thread همان consumer اجرا میشود. InvalidOperationException برای race condition نادری گرفته میشود که Task بین بررسی IsCanceled و اجرای آن cancel شود.
استفاده:
var pcQ = new PCQueue (2); // Maximum concurrency of 2
string result = await pcQ.Enqueue (() => "That was easy!");
...
در نتیجه تمام مزایای Task ــ propagation exception، return value و cancellation ــ را داریم و در عین حال scheduling را کاملاً خودمان کنترل میکنیم.