Страницы

Поиск по вопросам

Показаны сообщения с ярлыком tpl. Показать все сообщения
Показаны сообщения с ярлыком tpl. Показать все сообщения

четверг, 13 февраля 2020 г.

Количество потоков в Parallel.ForEach

#c_sharp #tpl


Устанавливаю количество потоков в 10, но оно работает в три потока (трехядерный комп). 

Parallel.ForEach(list, 
    new ParallelOptions { MaxDegreeOfParallelism = 10 },
    () => Thread.CurrentThread.ManagedThreadId,
    (i, loopState, threadId) =>
    {
        var cells = i.Split(';');
        var ip = cells[0];
        var login = cells[1];
        var pass = cells[2];
        try
        {
            using (var client = new SshClient(ip, login, pass))
            {
                client.Connect();
                var port = new ForwardedPortLocal("localhost", 10000 + (uint)threadId,
"ya.ru", 80);
                port.Start();
                client.AddForwardedPort(port);
                newList.Add(ip);
                port.Stop();
                client.Disconnect();
            }
        }
        catch 
        {

        }

        return threadId;
    },
    (threadId) => { }

    


Ответы

Ответ 1



В Parallel.ForEach количество потоков задается максимальным, а реальное количество потоков зависит от доступных ресурсов. Если нужно точное количество потоков, то можно использовать PLINQ либо использовать собственное решение с помощью кастомного TaskScheduler, либо явного контроля с помощью WaitHandle. Также не нужно забывать, что работа с сетью это так называемые "I/O bound потоки", поэтому лучше использовать такую асинхронность при условии, что используемая библиотека/класс это позволяет

Ответ 2



Насколько помню, у Parallel.ForEach довольно сложная логика, когда он сам подбирает наиболее оптимальное количество потоков для выполнения задачи. А MaxDegreeOfParallelism определяет лишь максимально возможное количество потоков. По своему опыту знаю, что количество активных потоков довольно сильно плавает, увеличиваясь и уменьшаясь в зависимости от характеристик компьютера, решаемых параллельно задач и того, насколько они занимают ресурсы. Хардкорный способ запустить сразу большое количество потоков - увеличить размер пула при старте приложения, например ThreadPool.SetMinThreads(200, 200); Parallel.ForEach(... Не буду говорить, что это хороший способ. Но тем не менее, мне когда такое решение помогло, когда надо было обрабатывать несколько сотен (или даже тысяч) параллельных потоков, выполняющих большие операции, связанные с вводом/выводом и сетевым взаимодействием.

Ответ 3



Если c Parallel.ForEach что-то не выходит - можно попытаться использовать задачи: Task.WaitAll(list.Select((item, index) => Task.StartNew(() => { /** тут код **/ /* здесь item - элемент списка, index - номер элемента в списке */ }))); Но если и с задачами не выйдет - проверьте загруженность процессора. SSH же использует криптографию, а криптография процессор неплохо нагружает...

Ответ 4



Как указывает название опции MaxDegreeOfParallelism -- это максимальное количество потоков. Опция является скорее подсказкой, нежели приказом. Фактическое же количество потоков подбирается библиотекой, исходя из количества ядер и загруженности системы. Стоит отметить, что потоки занимают память, а на переключение между ними требуется время. Оптимальное количество потоков в системе на самом деле равно количеству ядер -- все потоки будут заняты работой, отсутствуют траты на переключение контекста. Понятно, что в работающей ОС запущено множество потоков, поэтому библиотека оптимальный вариант -- запустить столько же потоков, сколько у вас ядер. Если вы хотите завести фиксированное количество потоков -- воспользуйтесь тасками. Алгоритм следующий: IEnumerable tasks = GetTasks(); int maxDegreeOfParallelism = 10; var queue = new Queue(maxDegreeOfParallelism); foreach (var task in tasks) { queue.Enqueue(task); if (queue.Count >= maxDegreeOfParallelism) { await queue.Dequeue(); } } while (queue.Count > 0) { await queue.Dequeue(); } Однако нужно тестировать: может случиться, что принудительный запуск 10 тасков (потоков) будет работать медленнее, чем 3 потока. Особенно в вашем случае, когда SSH использует процессор и сеть, а код у вас синхронный. Если же существует асинхронная версия SshClient, и вы правильным образом напишете весь код асинхронно, то теоретически получите выигрыш даже меньшим количеством потоков, поскольку пока данные ходят по интернету, поток может не ждать их, а заняться другой работой (например, обработкой уже пришедших результатов).

суббота, 8 февраля 2020 г.

Чем чревато отсутствие обработки OperationCanceledException у Task?

#c_sharp #многопоточность #tpl


Чем может быть чреват такой вот Task с необработанным исключением отмены действия,
если далее я к нему нигде не обращаюсь?

public static void Main()
{
    var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10));
    Task.Run(() => SomeMethod(cts.Token), cts.Token);
}

public static void SomeMethod(CancellationToken token)
{
    while (true)
    {
        token.ThrowIfCancellationRequested();
        // Некая трудоемкая операция
        Thread.Sleep(TimeSpan.FromMilliseconds(100));
    }
}


UPDATE1:
После ответа andreycha я решил проверить, как работает ловля ошибок у Task. И я заметил,
что исключение отмены не доходит до глобального обработчика и никак не вешает систему.
В кофиг файле прописал:  


  



Код:  

public static void Main()
{
    TaskScheduler.UnobservedTaskException += (object sender, UnobservedTaskExceptionEventArgs
eventArgs) =>
    {
        Console.WriteLine("Task error");
        eventArgs.SetObserved();
        (eventArgs.Exception).Handle(ex =>
        {
            Console.WriteLine("Exception type: " + ex.GetType());
            return true;
        });
    };

    CancellationTokenSource cts = new CancellationTokenSource(1000);
    Task.Factory.StartNew(() =>
    {
        Console.WriteLine("Enter");
        while (true)
        {
            cts.Token.ThrowIfCancellationRequested();
            Thread.Sleep(100);
        }
    }, cts.Token);
    Task.Factory.StartNew(() =>
    {
        throw new Exception("Some exception");
    });

    Thread.Sleep(4000);
    Console.WriteLine("Collecting");
    GC.Collect();
    GC.WaitForPendingFinalizers();

    Console.ReadLine();
}


Пример вывода:  

Task error
Exception type: System.Exception


Выходит, что если тебе никак не надо обработать отмену операции, то можно в обще
ничего не делать и ничего не будет все таки?
    


Ответы

Ответ 1



Если приложение работает под .NET Framework 4 (или под .NET Framework 4.5+ с включенной опцией ThrowUnobservedTaskExceptions), то когда сборщик мусора доберется до этого таска, финализатор выбросит исключение и приложение упадет. В .NET Framework 4.5+ поведение изменили и приложение продолжит работать (а само исключение по-прежнему можно отловить в обработчике UnobservedTaskException). Необработанное исключение ничему не помешает. Однако я бы рекомендовал всегда обзервить таски, иначе вы не узнаете, завершился ли таск успешно или нет и завершился ли вообще. В 99% случаев unobserved task -- это ошибка. Отследить завершение можно разными способами (зависит от вашей текущей архитектуры по большей степени): 1) Если вы уже используете async/await, тогда ожидайте так: var task = Task.Run(() => SomeMethod(cts.Token), cts.Token); try { await task; } catch (OperationCanceledException) { // задача была отменена } catch (Exception) { // другая ошибка } 2) Если ваш код полностью синхронный, то можно использовать либо продолжения, либо синхронное ожидание. Вариант с продолжением. Помните о том, что продолжение выполняется в том же контексте, что и оригинальный таск (т.е. в потоке из пула потоков). А значит обращаться напрямую к компонентам UI, например, нельзя. Task.Run(() => SomeMethod(cts.Token), cts.Token) .ContinueWith(SomeMethodHandler, TaskContinuationOption.OnlyOnFaulted); ... private void SomeMethodHandler(Task task) { if (task.Exception is OperationCanceledException) { // задача была отменена } else { // другая ошибка } } Вариант с синхронным ожиданием. Тут надо быть аккуратным с тем, в каком конкретно месте вы ожидаете. Поскольку внутри SomeMethod у вас бесконечный цикл, то на строке task.Wait() приложение будет висеть до тех пор, пока задача не будет отменена var task = Task.Run(() => SomeMethod(cts.Token), cts.Token); try { task.Wait(); } catch (AggregateException e) { // синхронное ожидание, в отличие от await, не "разворачивает" исключения // проверяем e.InnerExceptions на предмет наличия OperationCanceledException } P.S. Если же говорить о коде, который вы привели, то он завершит свое выполнение почти моментально, не успев произвести нужную работу. Потому что таск никто не ожидает.

Ответ 2



Если использовать библиотеку NLog - то там у класса Logger есть метод SwallowAsync, который позволяет залогировать асинхронную ошибку если она вдруг возникла. Впрочем, стандартная реализация тут в любом случае не подойдет, ведь отмена задачи - нормальный способ окончания работы, а не ошибочный. В данном случае, когда весь код свой, я бы предпочел не создавать себе проблем вместо того чтобы их героически решать: public static void SomeMethod(CancellationToken token) { while (true) { if (token.IsCancellationRequested) return; // Некая трудоемкая операция Thread.Sleep(TimeSpan.FromMilliseconds(100)); } }

четверг, 9 января 2020 г.

Task vs Thread: на каком ядре

#c_sharp #многопоточность #tpl #task


Про Task:


  Данная библиотека позволяет распараллелить задачи и выполнять их сразу
  на нескольких процессорах, если на целевом компьютере имеется
  несколько ядер. Кроме того, упрощается сама работа по созданию новых
  потоков. Поэтому начиная с .NET 4.0. рекомендуется использовать именно
  TPL и ее классы для создания многопоточных приложений, хотя
  стандартные средства и класс Thread по-прежнему находят широкое
  применение.


Источник: ссылка



Изначально кол-во потоков в пуле потоков равно числу ядер процессора. Вопрос: каждый
поток из потока пулов будет выполняться на "своем" ядре?


А как дела обстоят с обычными потоками - Thread (понимаю, что Task - более высокая
абстракция Thread)? Новый Thread может выполняться как на ядре, где выполняется главный
поток, так и на другом ядре, т.е. не гарантирует выполнение на новом ядре?
    


Ответы

Ответ 1



В цитате явно смесь процессоров и ядер. Но распараллеливаются обычно по ядрам, а эти ядра могут принадлежать разным процессорам. Изначально кол-во потоков в пуле потоков равно числу ядер процессора. на усмотрение библиотеки. Вопрос: каждый поток из потока пулов будет выполняться на "своем" ядре? можно сделать так, что бы каждый поток исполнялся только на своем ядре (осознанно привязав их), но так обычно не делают - планировщик ОС обычно несколько умнее и будет их разбрасывать по ядрам по своему разумению. А как дела обстоят с обычными потоками - Thread (понимаю, что Task - более высокая абстракция Thread)? Да, Task - абстракция, которая прячет от программиста детали. Просто таска, это функция, которая выполняется внутри потока. Синхронизация, очереди и получения результата прячутся библиотекой. Новый Thread может выполняться как на ядре, где выполняется главный поток, так и на новом? Там, где будет удобнее ОС разместить его. Более того, поток может "гулять" по ядрам.

понедельник, 23 декабря 2019 г.

Как выполнить цикл еще раз в Parallel.For?

#c_sharp #tpl


Допустим я в 100 потоков качаю картинки. Если какая либо итерация вызвало исключение
- как его повторить по новой?

Parallel.For(0, newLst.Count, new ParallelOptions { MaxDegreeOfParallelism = 100
}, (i) =>
      {
           DownloadImage(newLst[i]);
      });

    


Ответы

Ответ 1



Сделять DownloadImageAsync, который дергать не в Parallel.For, а просто ограничив число одновременно выполняемых задач, например с помощью такого кода: public static IEnumerable> ForEachAsync( this IEnumerable source, Func> selector, int degreeOfParallelism) { Contract.Requires(source != null); Contract.Requires(selector != null); // We need to know all the items in the source before starting tasks var tasks = source.ToList(); int completedTask = -1; // Creating an array of TaskCompletionSource that would holds // the results for each operations var taskCompletions = new TaskCompletionSource[tasks.Count]; for(int n = 0; n < taskCompletions.Length; n++) taskCompletions[n] = new TaskCompletionSource(); // Partitioner would do all grunt work for us and split // the source into appropriate number of chunks for parallel processing foreach (var partition in Partitioner.Create(tasks).GetPartitions(degreeOfParallelism)) { var p = partition; // Loosing sync context and starting asynchronous // computation for each partition Task.Run(async () => { while (p.MoveNext()) { var task = selector(p.Current); // Don't want to use empty catch . // This trick just swallows an exception await task.ContinueWith(_ => { }); int finishedTaskIndex = Interlocked.Increment(ref completedTask); taskCompletions[finishedTaskIndex].FromTask(task); } }); } return taskCompletions.Select(tcs => tcs.Task); } Теперь можно будет сделать так: var tasks = newLst.ForEachAsync(n => DownloadImageAsyncWithRetry(n, retryCount), 100).ToList() И логику ретраинга просто впихнуть в DownloadImageAsyncWithRetry: public Task DownloadImageAsyncWithRetry(input) { var tsk = DownloadImageAsync(input); // Тут все зависит от того, каким именно образом определяется неудача. // Если это исключение, то вешаем ContinueWith, если это код возврата, // то проверяем его и пробуем повторить запрос. } Да, Parallel.For в этом случае идея - не очень, поскольку он предназначен прежде всего для CPU Intensive операций, а здесь явно IO Intensive. В этом случае логично сделать сам метод асинхронным, который будет дергать асинхронный API для загрузки картинок. З.Ы. Код ForEachAsync взять отсюда.

четверг, 19 декабря 2019 г.

Применение Task.WhenAll

#c_sharp #многопоточность #tpl


Подскажите пожалуйста, в чем особенность использования Task.WhenAll ?

Полазив в msdn понял, что он создает новую задачу, по завершении указанных задач
в параметре, однако, по сути он ничего не создает, кроме некой ссылки типа Task, с
которой я не могу понять что делать дальше. Если здесь создается просто ссылка, которой
я должен присвоить в дальнейшем новый объект типа Task, то тогда проще вызвать (имхо) 

Task.WaitAll(t, t1);
Task t2 = new Task(...);


или может быть я в неправильном направлении думаю. Подскажите пожалуйста пример использования
этого метода. Спасибо.
    


Ответы

Ответ 1



Сначала я неправильно прочитал вопрос и ответил про WaitAll. Допустим есть несколько экземплров Task, каждый из которых выполняет действие и возвращает результат. Затем, все эти результаты надо как-то обработать. WhenAll создаёт таск, который заканчивается, когда заканчиваются все таски в него переданные. При этом результатом этого таска будет массив результатов тасков аргументов. Также можно проанализировать состояние переданных тасков, их исключения и так далее. static void Main(string[] args) { var t1 = new Task(DoSomething1); var t2 = new Task(DoSomething2); t1.Start(); t2.Start(); var t3 = Task.WhenAll(t1,t2); Console.WriteLine(t3.Result.Sum()); } private static int DoSomething2() { return 3; } private static int DoSomething1() { return 5; }

Ответ 2



Как это вы не можете понять, что делать с объектом типа Task? С любой задачей можно сделать три вещи: дождаться ее окончания асинхронно await Task.WhenAll(a1, a2); // К этому моменту a1 и a2 уже завершились дождаться ее окончания синхронно (не очень полезный вариант - привожу для полноты картины) Task.WhenAll(a1, a2).Wait(); // тоже самое, что и Task.WaitAll(a1, a2) сформировать продолжение var a3 = Task.WhenAll(a1, a2).ContinueWith(...);

понедельник, 16 декабря 2019 г.

Ожидание в асинхронности

#c_sharp #async_await #tpl


Есть абстрактный пример:

async void Do() 
{
    ...
    await DownloadSomething();
    // какой-то другой код, который выполнится позже 
    ...
}

void FuncMain() 
{
    Do();
    //какой-то код
} 


Когда начинается "долгая"  операция DownloadSomething, управление передаётся в FuncMain,
а после, когда загрузка закончится, продолжается код после DownloadSomething. 
Вопрос: где удерживается await DownloadSomething? Или удерживается в каком-то потоке
из пула? 
    


Ответы

Ответ 1



Смотря, что стоит в DownloadSomething(). Если это IO то поток IO. который фиксирован и независим от приложения. DB, SQL, Download и т.д. в ту же копилку. В момент передачи идут сохрание информации текущей конфигурации потока и как его возвращать. Текущий поток высвобождается и по окончанию происходит обратный вызов из пула нового/исходного потока. (ConfigureAwait(true/false)).

Ответ 2



Чтобы вопрос не оставался без ответа: Дело в том, что async-метод не является методом в обычном понимании этого слова. С точки зрения внешнего кода, его выполнение заканчивается практически сразу (с первым await'ом, который ожидает неокончившийся Task*). На время ожидания метод не выполняется нигде. Это не буддистский коан, а реальная подробность имплементации async/await. По существу await выполняется так: код просто подписывает на окончание выполнения Task'а метод специального скрытого объекта, и завершает выполнение. При окончании работы Task'а метод получает управление, и при помощи довольно простых трюков (наподобие goto в середину кода) возобновляет выполнение кода async-метода. Таким образом, во время await'а метод не выполняется ни в каком потоке. *или tasklike

суббота, 14 декабря 2019 г.

TaskScheduler и балансировка по ядрам для процессов

#c_sharp #tpl


Как известно, TaskScheduler в TPL раскидывает таски по ядрам (хоть и не гарантирует это).

Возьмем другой случай - порождаются много копий процессов, где внутри поток с многочисленными
Thread.Sleep. В таком варианте поток намертво прилипнет к какому то ядру.

Вопрос в том, если этот поток переделать на task-модель, то будут ли эти таски перемалываться
на разных ядрах или же будут тяготеть к одному и тому же ядру? Хоть это и таски, но
по факту один поток разбивается на таски, чтобы избежать sleep и TaskScheduler может
тяготеть переиспользовать этот же поток (других то в пуле нет)
    


Ответы

Ответ 1



Окей, вам нужно перебрасывать выполнение между ядрами. Это можно сделать вот как. Подсчитываем количество ядер. Это легко: Environment.ProcessorCount. Запускаем столько UI-потоков, сколько у нас ядер. Для этого берём код отсюда, и заимствуем из него класс DispatcherThread. Каждый из них представляет собой поток, в который можно переключиться при помощи await AsyncHelper.RedirectTo(t.Dispatcher); (оттуда же). Нам нужно разбросать эти потоки по ядрам. Это можно сделать как описано здесь. Теперь в нашей async-функции, если мы хотим поменять ядро, просто пишем currentCore = (currentCore + 1) % Environment.ProcessorCount; await AsyncHelper.RedirectTo(threads[currentCore].Dispatcher); Полный код: using System; using System.Collections; using System.Collections.Generic; using System.Diagnostics; using System.Linq; using System.Runtime.CompilerServices; using System.Text; using System.Threading; using System.Threading.Tasks; using System.Windows.Threading; namespace SO5 { class Program { static List threads; static int currentCoreNo = 0; static void Main(string[] args) { threads = Enumerable.Range(0, Environment.ProcessorCount) .Select(coreNo => new CoreAffineDispatcherThread(coreNo)) .ToList(); Run().Wait(); foreach (var t in threads) t.Dispose(); } static async Task Run() { for (int i = 0; i < 10; i++) { await Task.Delay(100); currentCoreNo = (currentCoreNo + 1) % Environment.ProcessorCount; await AsyncHelper.RedirectTo(threads[currentCoreNo].Dispatcher); var t = Thread.CurrentThread; Console.WriteLine($"Task reporting from thread {t.ManagedThreadId}," + $" thread pool: {t.IsThreadPoolThread}"); } } public class CoreAffineDispatcherThread : IDisposable { public Dispatcher Dispatcher { get; private set; } Thread thread; public CoreAffineDispatcherThread(int coreNumber) { using (var barrier = new AutoResetEvent(false)) { thread = new Thread(() => { Dispatcher = Dispatcher.CurrentDispatcher; barrier.Set(); Thread.BeginThreadAffinity(); #pragma warning disable 618 // The call to BeginThreadAffinity guarantees stable results // for GetCurrentThreadId, so we ignore the obsolete warning int osThreadId = AppDomain.GetCurrentThreadId(); #pragma warning restore 618 // Find the ProcessThread for this thread. ProcessThread thread = Process.GetCurrentProcess() .Threads.Cast() .Where(t => t.Id == osThreadId) .Single(); // Set the thread's processor affinity var cpuMask = 1 << coreNumber; thread.ProcessorAffinity = new IntPtr(cpuMask); Dispatcher.Run(); Thread.EndThreadAffinity(); }); thread.SetApartmentState(ApartmentState.STA); thread.Start(); barrier.WaitOne(); } } public void Dispose() { Dispatcher.InvokeShutdown(); if (thread != Thread.CurrentThread) thread.Join(); } } } static class AsyncHelper { public static DispatcherRedirector RedirectTo(Dispatcher d) { return new DispatcherRedirector(d); } } public struct DispatcherRedirector : INotifyCompletion { public DispatcherRedirector(Dispatcher dispatcher) { this.dispatcher = dispatcher; } #region awaiter public DispatcherRedirector GetAwaiter() { // combined awaiter and awaitable return this; } #endregion #region awaitable public bool IsCompleted { get { // true means execute continuation inline return dispatcher.CheckAccess(); } } public void OnCompleted(Action continuation) { dispatcher.BeginInvoke(continuation); } public void GetResult() { } #endregion Dispatcher dispatcher; } } При тестовом пробеге выдаёт: Task reporting from thread 10, thread pool: False Task reporting from thread 11, thread pool: False Task reporting from thread 12, thread pool: False Task reporting from thread 13, thread pool: False Task reporting from thread 14, thread pool: False Task reporting from thread 15, thread pool: False Task reporting from thread 16, thread pool: False Task reporting from thread 9, thread pool: False Task reporting from thread 10, thread pool: False Task reporting from thread 11, thread pool: False

суббота, 7 декабря 2019 г.

Как сделать асинхронный IEnumerable?

#c_sharp #tpl


Делаю поиск, по нескольким сразу местам.

На вернем уровне пока что-то типа:

  foreach (var item in Plugins.SelectMany(p => p.Search(query)))
    Items.Add(new ViewModel(item));


А в каждой реализации метода Search:

public override IEnumerable Search(string name)
{
  await some resource
  foreach (var element in networkResult)
    ...
    yield return result;
}


На деле, хочу параллельный доступ ко всем поискам, чтобы каждый элемент появлялся
в UI когда он готов, а не когда закончится всё целиком, как это сейчас работает. Что
именно тут в таски оборачивать - нету хороших идей. Снаружи вроде логичнее выглядит
Task, но реализовывать как правильно - не понимаю.
    


Ответы

Ответ 1



Смотрите, это не так сложно. Вы устанавливаете Ix, nuget-пакеты System.Interactive и System.Interactive.Async. У вас появляется интерфейс IAsyncEnumerable и вспомогательные классы. Пользоваться можно вот так: static class FileEx { public static IAsyncEnumerable ReadLinesAsync(string path) { // IAsyncEnumerable обладает только одной функцией - создать энумератор return AsyncEnumerable.CreateEnumerable(() => { var stream = File.OpenText(path); string current = null; // создаём энумератор при помощи готовой фабрики return AsyncEnumerable.CreateEnumerator( // у StreamReader.ReadLineAsync нет перегрузки с ct, пичалько moveNext: async ct => (current = await stream.ReadLineAsync()) != null, current: () => current, dispose: stream.Dispose); }); } } Смотрите, что тут происходит. Для начала, одна и та же последовательность может пробегаться разными кусками кода вперемежку, поэтому состояние текущего обхода мы держим в энумераторе. В принципе, нам нужно было бы завести отдельный класс для энумератора, и держать в нём свойства. Но мы пойдём более модным путём, и будем держать данные в замыкании. Мы открываем StreamReader, заводим переменную для текущей строки. Асинхронная функция MoveNext итератора получает следующую строку из StreamReader'а, и проверяет результат на null (null означает конец файла). Функция Current просто выдаёт текущую строку. А функция Dispose закрывает в конце поток. Теперь с этим можно работать: class Program { static async Task Main(string[] args) { var lines = FileEx.ReadLinesAsync("text.txt"); using (var en = lines.GetEnumerator()) { while (await en.MoveNext()) Console.WriteLine(en.Current); } } } Или просто await FileEx.ReadLinesAsync("text.txt") .ForEachAsync(s => Console.WriteLine(s)); В следующей версии C# планируется поддержка асинхронных энумераторов прямо в языке. С ней наш пример запишется так: static class FileEx { public static async IAsyncEnumerable ReadLinesAsync(string path) { using (var stream = File.OpenText(path)) { string current; while ((current = await stream.ReadLineAsync()) != null) yield return current; }; } } class Program { static async Task Main(string[] args) { foreach await (var s in FileEx.ReadLinesAsync("text.txt")) Console.WriteLine(s); } } Смотрите, для конкретно вашего случая (вот такой код) вам стоит разобрать задачу на составные части, поскольку асинхронных штук в ней много. Потом можно будет связать их вместо. Начнём с получения страниц и их разбора. Список хостов получить просто, тут не нужна асинхронность: var hosts = ConfigStorage.Plugins .Where(p => p.GetParser().GetType() == typeof(Parser)) .Select(p => p.GetSettings().MainUri); Теперь, нам нужно по хосту получить список HtmlNodeCollection. Это «длинная» задача, выносим её в таск: async Task GetHostMangasAsync(string name, Uri host, CookieClient client) { var searchHost = new Uri(host, "search?q=" + WebUtility.UrlEncode(name)); var page = await Task.Run(() => Page.GetPage(searchHost, client)); if (!page.HasContent) return null; return await Task.Run(() => { var document = new HtmlDocument(); document.LoadHtml(page.Content); return document.DocumentNode.SelectNodes("//div[@class='tile col-sm-6']"); }); } Теперь, нам нужно из неасинхронной последовательности host'ов и Task'а, который получает из каждого хоста HtmlNodeCollection, получить асинхронную последовательность. Такого метода в Ix из коробки я не нашёл, но его легко сколотить самому. Сделаем его обобщённым, вдруг ещё понадобится. Код практически ничем не отличается от примера с File.ReadLinesAsync. static class AsyncEnumerableExtensions { public static IAsyncEnumerable SelectAsync( this IEnumerable seq, Func> selector) { return AsyncEnumerable.CreateEnumerable(() => { IEnumerator seqEnum = seq.GetEnumerator(); R current = default; return AsyncEnumerable.CreateEnumerator( moveNext: async ct => { if (!seqEnum.MoveNext()) return false; current = await selector(seqEnum.Current); return true; }, current: () => current, dispose: seqEnum.Dispose); }); } } Вооружившись этим, мы можем написать такое: IAsyncEnumerable GetSearchPages(string name) { var hosts = ConfigStorage.Plugins .Where(p => p.GetParser().GetType() == typeof(Parser)) .Select(p => p.GetSettings().MainUri); var client = new CookieClient(); return hosts.SelectAsync(host => GetHostMangasAsync(name, host, client))) .Where(nc => nc != null); } Проверка на null нужна, потому что GetHostMangasAsync может вернуть null. Отлично, переходим дальше. Итак, у нас снова есть неасинхронная коллекция HtmlNodeCollection, из каждого элемента которой мы может вытащить при помощи асинхронной функции (т. к. у нас есть обращение к сети) экземпляр IManga. Пишем код: async Task GetMangaFromNode(Uri host, CookieClient client, HtmlNode manga) { // Это переводчик, идем дальше. if (manga.SelectSingleNode(".//i[@class='fa fa-user text-info']") != null) return null; var image = manga.SelectSingleNode(".//div[@class='img']//a//img"); var imageUri = image?.Attributes.Single(a => a.Name == "data-original").Value; var mangaNode = manga.SelectSingleNode(".//h3//a"); var mangaUri = mangaNode.Attributes.Single(a => a.Name == "href").Value; var mangaName = mangaNode.Attributes.Single(a => a.Name == "title").Value; if (!Uri.TryCreate(mangaUri, UriKind.Relative, out Uri test)) return null; var result = Mangas.Create(new Uri(host, mangaUri)); result.Name = WebUtility.HtmlDecode(mangaName); if (imageUri != null) result.Cover = await client.DownloadDataAsync(imageUri); return result; } Нам нужно теперь их соединить. Это несложно. Единственная проблема — в GetMangaFromNode тоже нужен CookieClient, а у нас от спрятан внутри GetSearchPages. Окей, будем передавать его снаружи. Затем, у нас из GetSearchPages возвращается только HtmlNodeCollection, а нужен ещё и host. Модифицируем GetSearchPages: будем возвращать пары из хоста и коллекции HtmlNode, и принимать на вход CookieClient: IAsyncEnumerable<(Uri host, HtmlNodeCollection nodes)> GetSearchPages( string name, CookieClient client) { var hosts = ConfigStorage.Plugins .Where(p => p.GetParser().GetType() == typeof(Parser)) .Select(p => p.GetSettings().MainUri); return hosts.SelectAsync( async host => (host, nodes: await GetHostMangasAsync(name, host, client))) .Where(pair => pair.nodes != null); } Ну и комбинируем. У нас каждая синхронная коллекция HtmlNode при помощи асинхронной функции даёт коллекцию экземпляров IManga. Это делается снова при помощи нашего SelectAsync: IAsyncEnumerable GetFromHostAndNodes( Uri host, HtmlNodeCollection nodes, CookieClient client) => nodes.SelectAsync(node => GetMangaFromNode(host, client, node)); Теперь можно складывать паззл: public IAsyncEnumerable Search(string name) { var client = new CookieClient(); return GetSearchPages(name, client) .SelectMany(pair => GetFromHostAndNodes(pair.host, pair.nodes, client)) .Where(m => m != null); } Всё!

четверг, 5 декабря 2019 г.

CancellationToken: почему структура?

#c_sharp #tpl #task


Почему CancellationToken реализован как структура?
Ведь структура является типом значения, как тогда реализован данный механизм?


static void Main(string[] args)
{
    CancellationTokenSource cancelTokenSource = new CancellationTokenSource();
    CancellationToken token = cancelTokenSource.Token;

    Task task1 = new Task(() => Factorial(5, token));
    task1.Start();

    cancelTokenSource.Cancel();
}

static void Factorial(int x, CancellationToken token)
{
    int result = 1;
    for (int i = 1; i <= x; i++)
    {
        if (token.IsCancellationRequested)
        {
            Console.WriteLine("Операция прервана токеном");
            return;
        }

        result *= i;
        Console.WriteLine("Факториал числа {0} равен {1}", i, result);
        Thread.Sleep(5000);
    }
}


Ведь структура является типом значения и создается копия объекта при передачи параметра
CancellationToken token
    


Ответы

Ответ 1



Почему CancellationToken реализован как структура? Для борьбы за эффективность. В большинстве случаев, да, CancellationToken вполне мог бы быть и классом, одна аллокация ничего не меняет, так как многопоточный код обычно некритичен к паре лишних мелких аллокаций. Но ведь CancellationToken задумывался как общий механизм отмены. И существуют случаи, в которых аллокации критичны для пользователей. Если бы CancellationToken был классом, то в этих случаях пользователям фреймворка проходилось бы пользоваться самописной структурой, и, хуже того, они не смогли бы пользоваться библиотечными методами. (Например, если пользователь хочет в метод, который ожидает CancellationToken, передать CancellationToken.None, потому что ему нужна скорость и не нужна отмена.) Так что разработчики решили облегчить нам, пользователям, жизнь, и сделали CancellationToken таки структурой. По поводу реализации механизма: да, разработчики нарушили семантику, и CancellationToken ведёт себя как типичный класс, а не как типичная структура. Они реализовали это таким образом. В CancellationToken есть единственное поле private CancellationTokenSource m_source; (ссылка на исходники). В проверке на равенство сравнивается m_source, и при копировании m_source копируется тоже, так что копия токена ведёт себя как оригинал, и тем самым неотличима от него. В частности, если оригинал токена отменён, то его копия — тоже.

суббота, 30 ноября 2019 г.

Принудительная отмена задачи

#c_sharp #net #многопоточность #tpl


Везде написано, что работа с задачами- это кооперативный процесс, т.е задача должна
сама корректно завершится при первой просьбе из внешнего кода.

Но, что делать если кто-то подводит?

Например, я делаю Cancel на токене и даю на завершение некоторое время, но задача
не завершается, а код должен двигаться дальше. Оставлять висеть задачу?

Читал, что есть Thread.Abort, но его не рекомендуют использовать.



Пример с выносом стороннего кода в отдельный процесс был приведен.

Хотелось бы увидеть еще какие-нибудь способы, еще как минимум решение задачи через
appDomain.
    


Ответы

Ответ 1



Ну вот вам пример реализации. Сразу предупреждаю, кода будет много. Возьмём в качестве основы вот такую ненадёжную функцию: class EvilComputation { static Random random = new Random(); public static async Task Compute( int numberOfSeconds, double x, CancellationToken ct) { bool wellBehaved = random.Next(2) == 0; var y = x * x; var delay = TimeSpan.FromSeconds(numberOfSeconds); await Task.Delay(delay, wellBehaved ? ct : CancellationToken.None); return y; } } Вы видим, что функция плохая: она может в зависимости от случайных условий не реагировать на отмену. Что делать в этом случае? Вынесем функцию в отдельный процесс. Этот процесс можно будет убить без особого вреда для исходного процесса. Для того, чтобы вызвать функцию в другом процессе, нужно передать данные о вызове функции туда. Для связи используем, например, анонимные пайпы (можно использовать по сути что угодно). Я основываю код на этом примере: How to: Use Anonymous Pipes for Local Interprocess Communication. Для передачи данных будем использовать стандартное бинарное форматирование, раз уж мы не пошли через WCF. Нам нужны DTO-объекты, которые будут перебрасываться между процессами. Их нужно использовать в двух процессах — главном и вспомогательном (назовём его плагином), поэтому для DTO-типов понадобится отдельная сборка. Заводим сборку OutProcCommonData, кладём в неё следующие классы: namespace OutProcCommonData { [Serializable] public class Command // общий класс-предок для посылаемой команды { } [Serializable] public class Evaluate : Command // команда на вычисление { public int NumberOfSecondsToProcess; public double X; } [Serializable] public class Cancel : Command // команда на отмену { } } Далее, возвращаемый результат: namespace OutProcCommonData { [Serializable] public class Response // общий класс-предок для возвращаемого результата { } [Serializable] public class Result : Response // готовый результат вычислений { public double Y; } [Serializable] public class Error : Response // ошибка с текстом { public string Text; } [Serializable] public class Cancelled : Response // подтверждение отмены { } } Далее, наш плагин. Это отдельное консольное приложение (хотя, если мы не хотим видеть консоль и отладочный вывод, можно сделать его неконсольным). Протокол общения таков. Главная программа посылает Evaluate, а после него, возможно, Cancel. Плагин возвращает Result в случае успешного вычисления, Cancelled в случае полученного сигнала отмены и успешно отменённого вычисления, и Error в случае ошибки (например, нарушения протокола коммуникации). Вот обвязочный код: class Plugin { static int Main(string[] args) { // нам должны быть переданы два аргумента: хендл входящего и исходящего пайпов if (args.Length != 2) { Console.Error.WriteLine("Shouldn't be started directly"); return 1; } return new Plugin().Run(args[0], args[1]).Result; } BinaryFormatter serializer = new BinaryFormatter(); // для сериализации async Task Run(string hIn, string hOut) { Console.WriteLine("[Plugin] Running"); // открывем переданные пайпы using (var inStream = new AnonymousPipeClientStream(PipeDirection.In, hIn)) using (var outStream = new AnonymousPipeClientStream(PipeDirection.Out, hOut)) { try { var cts = new CancellationTokenSource(); // токен для отмены Console.WriteLine("[Plugin] Reading args"); // пытаемся десериализовать аргументы var args = SafeGet(inStream); if (args == null) { Console.WriteLine("[Plugin] Didn't get args"); // отправляем ошибку, если не удалось serializer.Serialize( outStream, new OutProcCommonData.Error() { Text = "Unrecognized input" }); // и выходим return 3; } Console.WriteLine("[Plugin] Got args, start compute and waiting cancel"); // запускаем вычисление var computeTask = EvilComputation.Compute( args.NumberOfSecondsToProcess, args.X, cts.Token); // параллельно запускаем чтение возможной отмены var waitForCancelTask = Task.Run(() => (OutProcCommonData.Cancel)serializer.Deserialize(inStream)); // дожидаемся одного из двух var winner = await Task.WhenAny(computeTask, waitForCancelTask); // если первой пришла отмена... if (winner == waitForCancelTask) { Console.WriteLine("[Plugin] Got cancel, cancelling computation"); // просим вычисление завершиться cts.Cancel(); } // окончания вычисления всё равно нужно дождаться Console.WriteLine("[Plugin] Awaiting computation"); // если вычисление отменится, здесь будет исключение var result = await computeTask; Console.WriteLine("[Plugin] Sending back result"); // отсылаем результат в пайп serializer.Serialize( outStream, new OutProcCommonData.Result() { Y = result }); // нормальный выход return 0; } catch (OperationCanceledException) { // мы успешно отменили задание, рапортуем Console.WriteLine("[Plugin] Sending cancellation"); serializer.Serialize( outStream, new OutProcCommonData.Cancelled()); return 2; } catch (Exception ex) { // возникла непредвиденная ошибка, рапортуем Console.WriteLine($"[Plugin] Sending error {ex.Message}"); serializer.Serialize( outStream, new OutProcCommonData.Error() { Text = ex.Message }); return 3; } } } // ну и вспомогательная функция, которая пытается читать данные из пайпа T SafeGet(Stream s) where T : class { try { return (T)serializer.Deserialize(s); } catch { return null; } } } Я не отлавливаю ошибки при записи в пайп, добавьте сами по вкусу. Теперь, главная программа. Она будет у нас отдельно от плагина (то есть, у нас получаются три сборки). class Program { static void Main(string[] args) => new Program().Run().Wait(); async Task Run() { var cts = new CancellationTokenSource(); try { var y = await ComputeOutProc(2, cts.Token); Console.WriteLine($"[Main] Result: {y}"); } catch (TimeoutException) { Console.WriteLine("[Main] Timed out"); } catch (OperationCanceledException) { Console.WriteLine("[Main] Cancelled"); } } const int SecondsToSend = 3; const int TimeoutSeconds = 5; const int CancelSeconds = 2; BinaryFormatter serializer = new BinaryFormatter(); async Task ComputeOutProc(double x, CancellationToken ct) { Process plugin = null; bool pluginStarted = false; try { // создаём исходящий и входящий пайпы using (var commandStream = new AnonymousPipeServerStream( PipeDirection.Out, HandleInheritability.Inheritable)) using (var responseStream = new AnonymousPipeServerStream( PipeDirection.In, HandleInheritability.Inheritable)) { Console.WriteLine("[Main] Starting plugin"); plugin = new Process() { StartInfo = { FileName = "OutProcPlugin.exe", Arguments = commandStream.GetClientHandleAsString() + " " + responseStream.GetClientHandleAsString(), UseShellExecute = false } }; // запускаем плагин с параметрами plugin.Start(); pluginStarted = true; Console.WriteLine("[Main] Started plugin"); commandStream.DisposeLocalCopyOfClientHandle(); responseStream.DisposeLocalCopyOfClientHandle(); void Send(Command c) { serializer.Serialize(commandStream, c); commandStream.Flush(); } try { // отсылаем плагину команду на вычисление Console.WriteLine("[Main] Sending evaluate request"); Send(new OutProcCommonData.Evaluate() { NumberOfSecondsToProcess = SecondsToSend, X = x }); Task responseTask; bool readyInTime; bool cancellationSent = false; // внутри этого блока при отмене будем отсылать команду плагину using (ct.Register(() => { Send(new OutProcCommonData.Cancel()); Console.WriteLine("[Main] Requested cancellation"); cancellationSent = true; })) { Console.WriteLine("[Main] Starting getting response"); // ожидаем получение ответа responseTask = Task.Run(() => (Response)serializer.Deserialize(responseStream)); // или таймаута var timeoutTask = Task.Delay(TimeSpan.FromSeconds(TimeoutSeconds)); var winner = await Task.WhenAny(responseTask, timeoutTask); readyInTime = winner == responseTask; } // если наступил таймаут, просим процесс вежливо завершить вычисления if (!readyInTime) { if (!cancellationSent) { Console.WriteLine("[Main] Not ready in time, sending cancel"); Send(new OutProcCommonData.Cancel()); } else { Console.WriteLine("[Main] Not ready in time, cancel sent"); } // и ждём ещё немного, ну или прихода ответа var timeoutTask = Task.Delay(TimeSpan.FromSeconds(CancelSeconds)); await Task.WhenAny(responseTask, timeoutTask); } // если до сих пор ничего не пришло, плагин завис, убиваем его if (!responseTask.IsCompleted) { Console.WriteLine("[Main] No response, killing plugin"); plugin.Kill(); // это завершит ожидание с исключением, по идее // в ранних версиях .NET нужно было бы поймать // это исключение // и уходим с исключением-таймаутом ct.ThrowIfCancellationRequested(); throw new TimeoutException(); } // здесь мы уверены, что ожидание завершилось Console.WriteLine("[Main] Obtaining response"); var response = await responseTask; // тут может быть брошено исключение // если была затребована отмена, выходим ct.ThrowIfCancellationRequested(); // проверяем тип результата: switch (response) { case Result r: // нормальный результат, возвращаем его Console.WriteLine("[Main] Got result, returning"); return r.Y; case Cancelled _: // отмена не по ct = таймаут Console.WriteLine("[Main] Got cancellation"); throw new TimeoutException(); case Error err: // пришла ошибка, бросаем исключение // лучше, конечно, определить собственный тип здесь Console.WriteLine("[Main] Got error"); throw new Exception(err.Text); default: // сюда мы вообще не должны попасть, если плагин работает нормально Console.WriteLine("[Main] Unexpected error"); throw new Exception("Unexpected response type"); } } catch (IOException e) { Console.WriteLine("[Main] IO error occured"); throw new Exception("IO Error", e); } } } finally { if (pluginStarted) { plugin.WaitForExit(); plugin.Close(); } } } } Результат пробега: [Main] Starting plugin [Main] Started plugin [Main] Sending evaluate request [Main] Starting getting response [Plugin] Running [Plugin] Reading args [Plugin] Got args, start compute and waiting cancel [Plugin] Awaiting computation [Plugin] Sending back result [Main] Obtaining response [Main] Got result, returning [Main] Result: 4 Если поменять константу SecondsToSend на 10, чтобы был таймаут, получаем такой результат двух пробегов: Для штатного завершения: [Main] Starting plugin [Main] Started plugin [Main] Sending evaluate request [Main] Starting getting response [Plugin] Running [Plugin] Reading args [Plugin] Got args, start compute and waiting cancel [Main] Not ready in time, sending cancel [Plugin] Got cancel, cancelling computation [Plugin] Awaiting computation [Plugin] Sending cancellation [Main] Obtaining response [Main] Got cancellation [Main] Timed out Для принудительного завершения: [Main] Starting plugin [Main] Started plugin [Main] Sending evaluate request [Main] Starting getting response [Plugin] Running [Plugin] Reading args [Plugin] Got args, start compute and waiting cancel [Main] Not ready in time, sending cancel [Plugin] Got cancel, cancelling computation [Plugin] Awaiting computation [Main] No response, killing plugin [Main] Timed out Если добавить перед var y = await ComputeOutProc(2, cts.Token); преждевременную отмену: cts.CancelAfter(TimeSpan.FromSeconds(1)); получим такой результат: для штатного завершения [Main] Starting plugin [Main] Started plugin [Main] Sending evaluate request [Main] Starting getting response [Plugin] Running [Plugin] Reading args [Plugin] Got args, start compute and waiting cancel [Main] Requested cancellation [Plugin] Got cancel, cancelling computation [Plugin] Awaiting computation [Plugin] Sending cancellation [Main] Obtaining response [Main] Cancelled и для принудительного завершения [Main] Starting plugin [Main] Started plugin [Main] Sending evaluate request [Main] Starting getting response [Plugin] Running [Plugin] Reading args [Plugin] Got args, start compute and waiting cancel [Main] Requested cancellation [Plugin] Got cancel, cancelling computation [Plugin] Awaiting computation [Main] Not ready in time, cancel sent [Main] No response, killing plugin [Main] Cancelled Наверняка кое-где недостаточно контролируются ошибки, так что проверяйте, не нужно ли ловить какие-то ещё исключения. На эту заготовку можно добавлять сверху свою логику. Например, можно аналогично пулу потоков завести пул плагинов, и доставлять задания свободному в данный момент плагину.

Ответ 2



Нашел один способ принудительного завершения задачи без нарушения работы приложения. Данный способ позволяет завершить задачу, запущенную с параметром TaskCreationOptions.LongRunning, зная ID рабочего потока. Основан на вызове функции ExitThread в контексте целевого потока с помощью недокументированной функции RtlRemoteCall. Способ не работает, если поток бесконечно находится в состоянии ожидания, можно завершить только работающий поток. Т.е это не на 100% надежно, но, я полагаю, получше чем Thread.Abort. Если задача исполняет ваш собственный код, ID рабочего потока легко получить, вызывая из нее функцию GetCurrentThreadId и сохраняя результат в переменной. Если задача исполняет чужой код (например, подгружаемый из внешней DLL), можно его узнать только выборкой из потоков по времени старта задачи. Основной класс: using System; using System.Collections.Generic; using System.Text; using System.Threading; using System.Threading.Tasks; using System.Diagnostics; using System.Runtime.InteropServices; namespace TaskTest { public class TaskKiller { const int THREAD_ACCESS_TERMINATE = (0x0001); const int SYNCHRONIZE = (0x00100000); const int STANDARD_RIGHTS_REQUIRED = (0x000F0000); const int THREAD_ALL_ACCESS = (STANDARD_RIGHTS_REQUIRED | SYNCHRONIZE | 0xFFFF); [DllImport("kernel32.dll")] public static extern IntPtr OpenThread(uint dwDesiredAccess, bool bInheritHandle, uint dwThreadId); [DllImport("kernel32.dll", SetLastError = true)] [return: MarshalAs(UnmanagedType.Bool)] public static extern bool CloseHandle(IntPtr hObject); [DllImport("kernel32.dll")] public static extern uint GetCurrentThreadId(); [DllImport("kernel32.dll")] public static extern IntPtr GetCurrentProcess(); [DllImport("kernel32.dll")] public static extern IntPtr GetCurrentThread(); [DllImport("kernel32.dll", SetLastError = true)] public static extern IntPtr GetModuleHandle(string lpModuleName); [DllImport("kernel32", CharSet = CharSet.Ansi, ExactSpelling = true, SetLastError = true)] public static extern IntPtr GetProcAddress(IntPtr hModule, string procName); [DllImport("ntdll.dll", ExactSpelling = true, EntryPoint = "RtlRemoteCall")] static extern int RtlRemoteCall( IntPtr Process, IntPtr Thread, IntPtr CallSite, uint ArgumentCount, IntPtr Arguments, uint PassContext, uint AlreadySuspended ); /// /// Завершение потока с указанным ID /// /// 0 при успешном завершении, код NTSTATUS при ошибке public static int KillThreadById(uint threadid) { IntPtr hModule = (IntPtr)0; IntPtr hThread = (IntPtr)0; IntPtr unmanagedPointer = (IntPtr)0; try { /*Получение адреса функции ExitThread*/ hModule = GetModuleHandle(@"kernel32.dll"); IntPtr pProc = (IntPtr)0; pProc = GetProcAddress(hModule, "ExitThread"); /*Получение дескриптора потока с полным доступом*/ hThread = TaskKiller.OpenThread( (uint)(THREAD_ALL_ACCESS), false, (uint)threadid); IntPtr hProcess = GetCurrentProcess(); int[] args = new int[] { (int)0 };//массив аргументов для RtlRemoteCall unmanagedPointer = Marshal.AllocHGlobal(args.Length * sizeof(int));//выделение блока неуправлемой памяти Marshal.Copy(args, 0, unmanagedPointer, args.Length);//копирование массива в неуправляемую память /*Вызов ExitThread в контексте завершаемого потока*/ int result = RtlRemoteCall(hProcess, hThread, pProc, 1, unmanagedPointer, 0, 0); return result; } finally { // Clean up resources if(unmanagedPointer !=(IntPtr)0)Marshal.FreeHGlobal(unmanagedPointer); if (hThread != (IntPtr)0) TaskKiller.CloseHandle(hThread); if (hModule != (IntPtr)0) TaskKiller.CloseHandle(hModule); } } /// /// Получение ID всех потоков, стартовавших в указанном интервале времени /// public static List GetThreadsByStartTime(DateTime t1, DateTime t2) { List threads = new List(); Process pr=Process.GetCurrentProcess(); using (pr) { ProcessThreadCollection ths = pr.Threads; foreach (ProcessThread th in ths) { using (th) { if (th.TotalProcessorTime.TotalMilliseconds > 0) { if (DateTime.Compare(th.StartTime, t1) >= 0 && DateTime.Compare(th.StartTime, t2) <= 0) threads.Add((uint)th.Id); } } } } return threads; } } } Пример использования: using System; using System.Collections.Generic; using System.Linq; using System.Text; using System.Windows.Forms; using System.Threading; using System.Threading.Tasks; using System.Diagnostics; using System.Runtime.InteropServices; namespace TaskTest { public partial class Form1 : Form { public Form1() { InitializeComponent(); PrintThreads(); } int i = 0; uint threadid;//ID рабочего потока DateTime t;//время старта задачи Task t1=null;//задача void PrintThreads() { Process pr=Process.GetCurrentProcess(); using (pr) { ProcessThreadCollection ths = pr.Threads; StringBuilder b = new StringBuilder(300); b.AppendLine("Threads: " + ths.Count); foreach (ProcessThread th in ths) { using (th) { b.AppendLine(th.Id + " - " + th.ThreadState.ToString()+" - "+th.StartTime); } } textBox1.Text += b.ToString(); } } private void button1_Click(object sender, EventArgs e) { /*Делегат для задачи*/ Action action = () => { threadid = TaskKiller.GetCurrentThreadId(); //сохранить ID потока для последующего доступа while (true)// Just loop. { i++; } }; // Construct an unstarted task t1 = new Task(action,TaskCreationOptions.LongRunning); // Launch task t1.Start(); t = DateTime.Now;//сохранить время старта для последующего использования textBox1.Text = "Task started"+Environment.NewLine; PrintThreads(); } private void button2_Click(object sender, EventArgs e) { //Завершение потока, если известен его ID textBox1.Text = "-- Before terminating --" + Environment.NewLine; PrintThreads(); textBox1.Text += Environment.NewLine; int res=TaskKiller.KillThreadById(threadid); if (res != 0) { textBox1.Text += ("Error NTSTATUS=" + res.ToString("X")); } else { textBox1.Text += threadid.ToString() + " is terminated!"; } textBox1.Text += Environment.NewLine; textBox1.Text += "-- After terminating --" + Environment.NewLine; PrintThreads(); } private void bTerminate_Click(object sender, EventArgs e) { //Завершение потоков по времени старта textBox1.Text = "-- Before terminating --" + Environment.NewLine; PrintThreads(); textBox1.Text += "-----------------------"; textBox1.Text += Environment.NewLine; List threads=TaskKiller.GetThreadsByStartTime( t.Subtract(TimeSpan.FromSeconds(1)), t.Add(TimeSpan.FromSeconds(1)) ); foreach (uint id in threads) { TaskKiller.KillThreadById(id); textBox1.Text += id.ToString() + " is terminated!"; textBox1.Text += Environment.NewLine; } textBox1.Text += "-- After terminating --" + Environment.NewLine; PrintThreads(); textBox1.Text += "-----------------------"; } } }

В чем разница между Task и Thread и когда что лучше использовать?

#c_sharp #net #многопоточность #tpl


Вроде, они предоставляют схожий функционал.
    


Ответы

Ответ 1



Это совсем разные вещи. Thread представляет собой физический, системный поток выполнения (за исключением SQL Server под .NET 2.0, да). А Task — это штука, которая по сути перепрыгивает из потока в поток, а зачастую и вовсе не находится ни в каком потоке! В результате у вас может быть всего 10 активных потоков, но тысячи Task'ов. Например, если вы делаете await на операцию чтения из сети, то он время ожидания прихода ответа от сервера ваша асинхронная функция вовсе не занимает никакого потока, а существует в спящем виде как обыкновенный объект где-то в памяти. Когда ответ реально приходит, функция находит какой-то поток (при обычных условиях это главный поток, но может быть и какой-то посторонний, если вы попросите), и продолжает выполнение на нём дальше. Для текущей версии языка имеет смысл почти всегда предпочитать Task'и и избегать Thread'ов, они слишком низкоуровневые. Пользуйтесь Task'ами, они умеют намного больше. Мне, например, за последний год пришлось использовать Thread только один раз (вот код), да и то в качестве дополнительного «хоста» для Task'ов. Необходимость была обусловлена тем, что мне нужен был STA apartment, а потоки из пула таковым не обладают.

Ответ 2



Я вырос из мира микроконтроллеров. И в этом мире была такая штука как кооперативная ОС. Так вот эти task и await очень похожи на кооперативную ОС.

В чем смысл TaskCompletionSource и когда его лучше использовать?

#c_sharp #net #асинхронность #tpl


Немного не понял смысла класса TaskCompletionSource.
В некоторых источниках пишут, что лучше его возвращать из метода вместо обычного
Task.Run().

Разве есть какой-то смысл? Что так, что так я смогу вызвать await на вызывающей стороне.
    


Ответы

Ответ 1



TaskCompletionSource — это тот самый крайний случай, когда вы не можете создать «базовый» Task стандартными средствами. Давайте я поясню, что я имею в виду. Если вы создаёте Task, обычных путей для этого два. Во-первых, если ваш код не производит ожидания, а активно работает, например, проводит вычисления (CPU-bound), вы отправляете его на пул потоков при помощи Task.Run или его аналогов. Во-вторых, если вы пользуетесь другими асинхронными операциями, вы создаёте async-метод, в котором производите await на другие асинхронные операции. .NET предоставляет множество готовых асинхронных операций, например, NetworkStream.ReadAsync или там Dispatcher.InvokeAsync. Но что делать, если вам нужно самому создать примитивную асинхронную операцию, которая не выражается в терминах других, уже готовых асинхронных операций? Как созданы самые внутренние Task-методы? В этом месте вам как раз и пригодится TaskCompletionSource. Например, вы хотите асинхронно дождаться события. Для этого вам нужно превратить событие в Task. Это делается как-то так: мы подписываемся на событие, и по его приходу завершаем Task. Task WaitInput() { var tcs = new TaskCompletionSource(); source.InputReceived += (o, args) => tcs.SetResult(args.Input); return tcs.Task; } Более строгий вариант с отпиской, в которой TaskCompletionSource используется как внутренний Task, чтобы успеть отписаться после его окончания: async Task WaitInput() { var tcs = new TaskCompletionSource(); SourceInputHandler handler = (o, args) => tcs.SetResult(args.Input); source.InputReceived += handler; try { return await tcs.Task; } finally { source.InputReceived -= handler; } } Ещё один пример из реального кода: запустить и дождаться окончания процесса: Task ExecuteProcess(string path) { var p = new Process() { EnableRaisingEvents = true, StartInfo = { FileName = path } }; var tcs = new TaskCompletionSource(); p.Exited += (sender, args) => { tcs.SetResult(true); p.Dispose(); }; // запуск выгружаем на пул потоков, потому что он медленный Task.Run(() => p.Start()); return tcs.Task; } Ещё один пример взят из класса DispatcherThread. Нам нужно дождаться, пока поток стартует, и придёт в «рабочее» состояние. Обычно для этого используют AutoResetEvent, но блокироваться в ожидании его неохота, и намного проще использовать TaskCompletionSource: static public Task CreateAsync() { var waitCompletionSource = new TaskCompletionSource(); var thread = new Thread(() => { // тут могут быть любые настройки waitCompletionSource.SetResult(new DispatcherThread()); Dispatcher.Run(); }); thread.SetApartmentState(ApartmentState.STA); thread.Start(); return waitCompletionSource.Task; } Мы видим, что таким образом можно превратить в Task по сути любую операцию. Дополнительное чтение по теме: TPL and Traditional .NET Framework Asynchronous Programming. Резюме: TaskCompletionSource позволяет превратить в Task даже ту асинхронную операцию, которая не даёт await-абельного интерфейса.

Ответ 2



На SO был такой вопрос. Почитай здесь. Перевод одного из лучших ответов: В моих опытах TaskCompletionSource отлично подходит для переноса старых асинхронных шаблонов в современный шаблон async/await. Самый полезный пример, о котором я могу думать, - это работать с Socket. Он имеет старые шаблоны APM и EAP, но не методы awaitable Task, которые имеют TcpListener и TcpClient. У меня лично есть несколько проблем с классом NetworkStream и предпочитаю raw Socket. Будучи тем, что мне также нравится шаблон async/await, я создал класс расширения SocketExtender, который создает несколько методов расширения для Socket. Все эти методы используют TaskCompletionSource для обертывания асинхронных вызовов следующим образом: public static Task AcceptAsync(this Socket socket) { if (socket == null) throw new ArgumentNullException("socket"); var tcs = new TaskCompletionSource(); socket.BeginAccept(asyncResult => { try { var s = asyncResult.AsyncState as Socket; var client = s.EndAccept(asyncResult); tcs.SetResult(client); } catch (Exception ex) { tcs.SetException(ex); } }, socket); return tcs.Task; } Я передаю Socket в методы BeginAccept, так что я получаю небольшое повышение производительности от компилятора, не требуя поднять локальный параметр. Тогда красота всего этого: var listener = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); listener.Bind(new IPEndPoint(IPAddress.Loopback, 2610)); listener.Listen(10); var client = await listener.AcceptAsync();

понедельник, 25 ноября 2019 г.

Что такое Task.Yield()?


Я не понимаю что это, как работает и в каких случаях используется. Может кто-нибудь по-русски объяснить?
    


Ответы

Ответ 1



Этот метод возвращает специальное значение, предназначенное для передачи оператору await, и в отрыве от этого оператора не имеющее смысла. Конструкция же await Task.Yield() делает довольно простую вещь — прерывает текущий метод и сразу же планирует его продолжение в текущем контексте синхронизации. Используется же эта конструкция для разных целей. Во-первых, эта конструкция может быть использована для немедленного возврата управлени вызывающему коду. Например, при вызове из обработчика события событие будет считаться обработанным: protected override async void OnClosing(CancelEventArgs e) { e.Cancel = true; await Task.Yield(); // (какая-то логика) } Во-вторых, эта конструкция используется для очистки синхронного контекста вызова. Например, так можно "закрыть" текущую транзакцию (ambient transaction): using (var ts = new TransactionScope()) { // ... Foo(); // ... ts.Complete(); } async void Foo() { // ... тут мы находимся в контексте транзакции if (Transaction.Current != null) await Task.Yield(); // ... а тут его уже нет! } В-третьих, эта конструкция может очистить стек вызовов. Это может быть полезным если программа падает с переполнением стека при обработке кучи вложенных продолжений. Например, рассмотрим упрощенную реализацию AsyncLock: class AsyncLock { private Task unlockedTask = Task.CompletedTask; public async Task Lock() { var tcs = new TaskCompletionSource(); await Interlocked.Exchange(ref unlockedTask, tcs.Task); return () => tcs.SetResult(null); } } Здесь поступающие запросы на получение блокировки выстраиваются в неявную очередь на продолжениях. Казалось бы, что может пойти не так? private static async Task Foo() { var _lock = new AsyncLock(); var unlock = await _lock.Lock(); for (var i = 0; i < 100000; i++) Bar(_lock); unlock(); } private static async void Bar(AsyncLock _lock) { var unlock = await _lock.Lock(); // do something sync unlock(); } Здесь продолжение метода Bar вызывается в тот момент, когда другой метод Bar выполняе вызов unlock(). Получается косвенная рекурсия между методом Bar и делегатом unlock, которая быстро сжирает стек и ведет к его переполнению. Добавление же вызова Task.Yield() перенесет исполнение в "чистый" фрейм стека, ошибка исчезнет: class AsyncLock { private Task unlockedTask = Task.CompletedTask; public async Task Lock() { var tcs = new TaskCompletionSource(); var prevTask = Interlocked.Exchange(ref unlockedTask, tcs.Task); if (!prevTask.IsCompleted) { await prevTask; await Task.Yield(); } return () => tcs.SetResult(null); } } Кстати, альтернативный способ починить код выше — использование флага RunContinuationsAsynchronously: class AsyncLock { private Task unlockedTask = Task.CompletedTask; public async Task Lock() { var tcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); await Interlocked.Exchange(ref unlockedTask, tcs.Task); return () => tcs.SetResult(null); } } В-четвертых, при использовании в UI-потоке эта конструкция позволяет обработать накопившиеся события ввода-вывода, что полезно при длительных обновлениях интерфейса. Например, при добавлении миллиона строк в таблицу программа не будет реагироват на действия пользователя, пока все строки не будут добавлены. Но если, к примеру, после добавления каждой тысячи строк вставлять вызов await Task.Yield() - программа сможет обрабатывать действия пользователя и не будет выглядеть зависшей. В WinForms для тех же целей можно было использовать метод Application.DoEvents( - но его избыточное использование приводило к переполнению стека. await Task.Yield() - это универсальный способ, который можно использовать как в WinForms, так и в WPF.

Ответ 2



Я думаю, что здесь никто не ответил на вопрос, зачем нужна Task.Yield. Она нужна когда задача (Task) использует бесконечный цикл (вообще любая продолжительная синхронная работа) и может удерживать поток из пула только для себя и не давать други задачам использовать этот поток. Task.Yield переотправляет задачу в очередь пула потоков и другие задачи, которые ожидали выполнения смогут использовать данный удерживаемый поток. Пример: CancellationTokenSource cts; void Start() { cts = new CancellationTokenSource(); // Запускаем асинхронную операцию var task = Task.Run(() => SomeWork(cts.Token), cts.Token); // Ждем окончания // После окончания операции обрабатываем результат/отмену/исключения } async Task SomeWork(CancellationToken cancellationToken) { int result = 0; bool loopAgain = true; while (loopAgain) { // Что-то делаем ... loopAgain = /* проверка на окончание цикла && */ cancellationToken.IsCancellationRequested; if (loopAgain) { // переотправляет задачу в очередь пула потоков чтобы другие задачи, которые ожидали выполнения смогли использовать данный поток await Task.Yield(); } } cancellationToken.ThrowIfCancellationRequested(); return result; } void Cancel() { // Запрашиваем отмену операции cts.Cancel(); }

среда, 10 июля 2019 г.

Как передать коллекцию прямоугольников в ItemsControl с Canvas асинхронно?

При решении вопроса, возник новый.
Что делаю: из ViewModel передаю коллекцию прямоугольников, вот так:
public async void Start() { RectItems.Clear();
CrossStitch cs = new CrossStitch() { BlockSize = _blockSize, Source = _sourceImage };
var data = await cs.Create();
foreach (var r in data) RectItems.Add(r);
}
Получаю во View вот так:

Все замечательно работало, до того как метод Start() стал async. Теперь я получаю вместо результата, это:
Необходимо создать DependencySource в том же потоке, в котором создан DependencyObject.
Нашел только одну похожую проблему, но в ней передавалось изображение в Canvas, и проблема решалась вызовом Freeze у изображения. А как быть в моем случае?
Метод Create и прилежащие:
public Task> Create() { return PixelateAsync(_source); }
private Task> PixelateTask(Bitmap source) { return Task.Factory.StartNew(() => Pixelate(source)); }
private Task> PixelateAsync(Bitmap source) { return PixelateTask(source); }
Pixelate(source) синхронный.
Класс RectItem
public class RectItem { public double X { get; set; } public double Y { get; set; } public double Width { get; set; } public double Height { get; set; } public System.Windows.Media.Brush C { get; set; } }
Метод Pixelate весь:
private List Pixelate(Bitmap source) { var result = new Bitmap(source);
List rectangs = new List();
using (var graphics = Graphics.FromImage(result)) { graphics.PageUnit = GraphicsUnit.Pixel;
for (int x = 0; x < source.Width; x += _blockSize) { for (int y = 0; y < source.Height; y += _blockSize) { var sums = new Sums();
for (int xx = 0; xx < _blockSize; ++xx) { for (int yy = 0; yy < _blockSize; ++yy) { if (x + xx >= source.Width || y + yy >= source.Height) { continue; }
var color = source.GetPixel(x + xx, y + yy); sums.A += color.A; sums.R += color.R; sums.G += color.G; sums.B += color.B; sums.T++; } }
var average = Color.FromArgb( sums.A / sums.T, sums.R / sums.T, sums.G / sums.T, sums.B / sums.T);
average = GetNearestColor(average); System.Windows.Media.Color mcolor = System.Windows.Media.Color.FromArgb(average.A, average.R, average.G, average.B); System.Windows.Media.Brush brush = new System.Windows.Media.SolidColorBrush(mcolor); rectangs.Add(new RectItem() { X = x + BlockSize, Y = y + BlockSize, Height = BlockSize, Width = BlockSize, C = brush }); } } }
return rectangs; }


Ответ

Смотрите. Проблема в том, что VM-классы создаются в фоновом потоке, это в обычной ситуации неправильно.
Но в вашем случае RectItem — не DependencyObject, а значит, он не привязан к определённому потоку. Поэтому можно пойти более простым путём: создавать этот объект где угодно. Единственная проблема, которую нужно вынести в UI — создание Brush. Но Brush является Freezable- значит, его можно также создавать где угодно, просто нужно после создания вызвать brush.Freeze();
Ещё один framework-класс — Color — тоже не является проблемой, т. к. он не является ни DependencyObject'ом, ни Freezable
Итого: просто добавьте после
System.Windows.Media.Brush brush = new System.Windows.Media.SolidColorBrush(mcolor);
строку
brush.Freeze();

пятница, 14 июня 2019 г.

Условия на тип возвращаемого значения метода при использовании await?

Есть метод:
public async T Method() { T result = await doSomeStuff();
return result; }
Какие условия должны быть выполнены для T, чтобы этот метод можно было вызвать:
public async void AnotherMethod() { await Method(); }


Ответ

По идее, await можно использовать для любого типа, в котором есть метод GetAwaiter, возвращающий реализацию интерфейса INotifyCompletion
public AlexsAwaiter GetAwaiter() { return new AlexsAwaiter(); }
class AlexsAwaiter : INotifyCompletion { public bool IsCompleted { get { ... } } public void OnCompleted(Action continuation) { ... } public void GetResult() { ... } }
Тип, возвращаемый GetResult и будет возвращаться из await
Подробнее лучше почитать например в книге
Дэвис Д. - Асинхронное программирование в C# 5.0 - 2013г.
или тут

среда, 12 июня 2019 г.

Task.IsComplited до реального завершения задачи

Создаю кучу Task - в каждом игровой цикл, помещаю их в List
GamesList.Add(gp.ContinueWith(t=>GamesList.Remove(t)));
Но они почему-то удаляются из списка до того как игровой цикл завершится - игра при этом идет спокойненько. Wtf?


Ответ

Проблема в том, что вы добавляете в список не тот таск, который вы удаляете.
ContinueWith возвращает новый таск, а пытаетесь удалять вы первоначальный таск.
И да, вы не должны работать со списком из разных потоков без блокировки. Ну то есть вы можете, но не удивляйтесь тогда потерянным данным.

Попробуйте заменить GamesList на потокобезопасную коллекцию, и используйте
GamesList.Add(gp); gp.ContinueWith(t=>GamesList.Remove(t));

пятница, 12 апреля 2019 г.

Чем чревато отсутствие обработки OperationCanceledException у Task?

Чем может быть чреват такой вот Task с необработанным исключением отмены действия, если далее я к нему нигде не обращаюсь?
public static void Main() { var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); Task.Run(() => SomeMethod(cts.Token), cts.Token); }
public static void SomeMethod(CancellationToken token) { while (true) { token.ThrowIfCancellationRequested(); // Некая трудоемкая операция Thread.Sleep(TimeSpan.FromMilliseconds(100)); } }
UPDATE1: После ответа andreycha я решил проверить, как работает ловля ошибок у Task. И я заметил, что исключение отмены не доходит до глобального обработчика и никак не вешает систему. В кофиг файле прописал:

Код:
public static void Main() { TaskScheduler.UnobservedTaskException += (object sender, UnobservedTaskExceptionEventArgs eventArgs) => { Console.WriteLine("Task error"); eventArgs.SetObserved(); (eventArgs.Exception).Handle(ex => { Console.WriteLine("Exception type: " + ex.GetType()); return true; }); };
CancellationTokenSource cts = new CancellationTokenSource(1000); Task.Factory.StartNew(() => { Console.WriteLine("Enter"); while (true) { cts.Token.ThrowIfCancellationRequested(); Thread.Sleep(100); } }, cts.Token); Task.Factory.StartNew(() => { throw new Exception("Some exception"); });
Thread.Sleep(4000); Console.WriteLine("Collecting"); GC.Collect(); GC.WaitForPendingFinalizers();
Console.ReadLine(); }
Пример вывода:
Task error Exception type: System.Exception
Выходит, что если тебе никак не надо обработать отмену операции, то можно в обще ничего не делать и ничего не будет все таки?


Ответ

Если приложение работает под .NET Framework 4 (или под .NET Framework 4.5+ с включенной опцией ThrowUnobservedTaskExceptions), то когда сборщик мусора доберется до этого таска, финализатор выбросит исключение и приложение упадет.
В .NET Framework 4.5+ поведение изменили и приложение продолжит работать (а само исключение по-прежнему можно отловить в обработчике UnobservedTaskException). Необработанное исключение ничему не помешает.

Однако я бы рекомендовал всегда обзервить таски, иначе вы не узнаете, завершился ли таск успешно или нет и завершился ли вообще. В 99% случаев unobserved task -- это ошибка. Отследить завершение можно разными способами (зависит от вашей текущей архитектуры по большей степени):
1) Если вы уже используете async/await, тогда ожидайте так:
var task = Task.Run(() => SomeMethod(cts.Token), cts.Token); try { await task; } catch (OperationCanceledException) { // задача была отменена } catch (Exception) { // другая ошибка }
2) Если ваш код полностью синхронный, то можно использовать либо продолжения, либо синхронное ожидание.
Вариант с продолжением. Помните о том, что продолжение выполняется в том же контексте, что и оригинальный таск (т.е. в потоке из пула потоков). А значит обращаться напрямую к компонентам UI, например, нельзя.
Task.Run(() => SomeMethod(cts.Token), cts.Token) .ContinueWith(SomeMethodHandler, TaskContinuationOption.OnlyOnFaulted); ... private void SomeMethodHandler(Task task) { if (task.Exception is OperationCanceledException) { // задача была отменена } else { // другая ошибка } }
Вариант с синхронным ожиданием. Тут надо быть аккуратным с тем, в каком конкретно месте вы ожидаете. Поскольку внутри SomeMethod у вас бесконечный цикл, то на строке task.Wait() приложение будет висеть до тех пор, пока задача не будет отменена
var task = Task.Run(() => SomeMethod(cts.Token), cts.Token); try { task.Wait(); } catch (AggregateException e) { // синхронное ожидание, в отличие от await, не "разворачивает" исключения // проверяем e.InnerExceptions на предмет наличия OperationCanceledException }

P.S. Если же говорить о коде, который вы привели, то он завершит свое выполнение почти моментально, не успев произвести нужную работу. Потому что таск никто не ожидает.

понедельник, 18 февраля 2019 г.

Task vs Thread: на каком ядре

Про Task:
Данная библиотека позволяет распараллелить задачи и выполнять их сразу на нескольких процессорах, если на целевом компьютере имеется несколько ядер. Кроме того, упрощается сама работа по созданию новых потоков. Поэтому начиная с .NET 4.0. рекомендуется использовать именно TPL и ее классы для создания многопоточных приложений, хотя стандартные средства и класс Thread по-прежнему находят широкое применение.
Источник: ссылка

Изначально кол-во потоков в пуле потоков равно числу ядер процессора. Вопрос: каждый поток из потока пулов будет выполняться на "своем" ядре? А как дела обстоят с обычными потоками - Thread (понимаю, что Task - более высокая абстракция Thread)? Новый Thread может выполняться как на ядре, где выполняется главный поток, так и на другом ядре, т.е. не гарантирует выполнение на новом ядре?


Ответ

В цитате явно смесь процессоров и ядер. Но распараллеливаются обычно по ядрам, а эти ядра могут принадлежать разным процессорам.
Изначально кол-во потоков в пуле потоков равно числу ядер процессора.
на усмотрение библиотеки.
Вопрос: каждый поток из потока пулов будет выполняться на "своем" ядре?
можно сделать так, что бы каждый поток исполнялся только на своем ядре (осознанно привязав их), но так обычно не делают - планировщик ОС обычно несколько умнее и будет их разбрасывать по ядрам по своему разумению.
А как дела обстоят с обычными потоками - Thread (понимаю, что Task - более высокая абстракция Thread)?
Да, Task - абстракция, которая прячет от программиста детали. Просто таска, это функция, которая выполняется внутри потока. Синхронизация, очереди и получения результата прячутся библиотекой.
Новый Thread может выполняться как на ядре, где выполняется главный поток, так и на новом?
Там, где будет удобнее ОС разместить его. Более того, поток может "гулять" по ядрам.

четверг, 1 ноября 2018 г.

Применение Task.WhenAll

Подскажите пожалуйста, в чем особенность использования Task.WhenAll ?
Полазив в msdn понял, что он создает новую задачу, по завершении указанных задач в параметре, однако, по сути он ничего не создает, кроме некой ссылки типа Task, с которой я не могу понять что делать дальше. Если здесь создается просто ссылка, которой я должен присвоить в дальнейшем новый объект типа Task, то тогда проще вызвать (имхо)
Task.WaitAll(t, t1); Task t2 = new Task(...);
или может быть я в неправильном направлении думаю. Подскажите пожалуйста пример использования этого метода. Спасибо.


Ответ

Сначала я неправильно прочитал вопрос и ответил про WaitAll
Допустим есть несколько экземплров Task, каждый из которых выполняет действие и возвращает результат. Затем, все эти результаты надо как-то обработать. WhenAll создаёт таск, который заканчивается, когда заканчиваются все таски в него переданные. При этом результатом этого таска будет массив результатов тасков аргументов.
Также можно проанализировать состояние переданных тасков, их исключения и так далее.
static void Main(string[] args) { var t1 = new Task(DoSomething1); var t2 = new Task(DoSomething2);
t1.Start(); t2.Start();
var t3 = Task.WhenAll(t1,t2);
Console.WriteLine(t3.Result.Sum()); }
private static int DoSomething2() { return 3; }
private static int DoSomething1() { return 5; }

среда, 24 октября 2018 г.

TaskScheduler и балансировка по ядрам для процессов

Как известно, TaskScheduler в TPL раскидывает таски по ядрам (хоть и не гарантирует это).
Возьмем другой случай - порождаются много копий процессов, где внутри поток с многочисленными Thread.Sleep. В таком варианте поток намертво прилипнет к какому то ядру.
Вопрос в том, если этот поток переделать на task-модель, то будут ли эти таски перемалываться на разных ядрах или же будут тяготеть к одному и тому же ядру? Хоть это и таски, но по факту один поток разбивается на таски, чтобы избежать sleep и TaskScheduler может тяготеть переиспользовать этот же поток (других то в пуле нет)


Ответ

Окей, вам нужно перебрасывать выполнение между ядрами. Это можно сделать вот как.
Подсчитываем количество ядер. Это легко: Environment.ProcessorCount Запускаем столько UI-потоков, сколько у нас ядер. Для этого берём код отсюда, и заимствуем из него класс DispatcherThread. Каждый из них представляет собой поток, в который можно переключиться при помощи await AsyncHelper.RedirectTo(t.Dispatcher); (оттуда же). Нам нужно разбросать эти потоки по ядрам. Это можно сделать как описано здесь Теперь в нашей async-функции, если мы хотим поменять ядро, просто пишем
currentCore = (currentCore + 1) % Environment.ProcessorCount; await AsyncHelper.RedirectTo(threads[currentCore].Dispatcher);

Полный код:
using System; using System.Collections; using System.Collections.Generic; using System.Diagnostics; using System.Linq; using System.Runtime.CompilerServices; using System.Text; using System.Threading; using System.Threading.Tasks; using System.Windows.Threading;
namespace SO5 { class Program { static List threads; static int currentCoreNo = 0;
static void Main(string[] args) { threads = Enumerable.Range(0, Environment.ProcessorCount) .Select(coreNo => new CoreAffineDispatcherThread(coreNo)) .ToList(); Run().Wait();
foreach (var t in threads) t.Dispose(); }
static async Task Run() { for (int i = 0; i < 10; i++) { await Task.Delay(100);
currentCoreNo = (currentCoreNo + 1) % Environment.ProcessorCount; await AsyncHelper.RedirectTo(threads[currentCoreNo].Dispatcher);
var t = Thread.CurrentThread; Console.WriteLine($"Task reporting from thread {t.ManagedThreadId}," + $" thread pool: {t.IsThreadPoolThread}"); } }
public class CoreAffineDispatcherThread : IDisposable { public Dispatcher Dispatcher { get; private set; }
Thread thread;
public CoreAffineDispatcherThread(int coreNumber) { using (var barrier = new AutoResetEvent(false)) { thread = new Thread(() => { Dispatcher = Dispatcher.CurrentDispatcher; barrier.Set(); Thread.BeginThreadAffinity();
#pragma warning disable 618 // The call to BeginThreadAffinity guarantees stable results // for GetCurrentThreadId, so we ignore the obsolete warning int osThreadId = AppDomain.GetCurrentThreadId(); #pragma warning restore 618
// Find the ProcessThread for this thread. ProcessThread thread = Process.GetCurrentProcess() .Threads.Cast() .Where(t => t.Id == osThreadId) .Single(); // Set the thread's processor affinity var cpuMask = 1 << coreNumber; thread.ProcessorAffinity = new IntPtr(cpuMask);
Dispatcher.Run();
Thread.EndThreadAffinity(); });
thread.SetApartmentState(ApartmentState.STA); thread.Start(); barrier.WaitOne(); } }
public void Dispose() { Dispatcher.InvokeShutdown(); if (thread != Thread.CurrentThread) thread.Join(); } } }
static class AsyncHelper { public static DispatcherRedirector RedirectTo(Dispatcher d) { return new DispatcherRedirector(d); } }
public struct DispatcherRedirector : INotifyCompletion { public DispatcherRedirector(Dispatcher dispatcher) { this.dispatcher = dispatcher; }
#region awaiter public DispatcherRedirector GetAwaiter() { // combined awaiter and awaitable return this; } #endregion
#region awaitable public bool IsCompleted { get { // true means execute continuation inline return dispatcher.CheckAccess(); } }
public void OnCompleted(Action continuation) { dispatcher.BeginInvoke(continuation); }
public void GetResult() { } #endregion
Dispatcher dispatcher; } }
При тестовом пробеге выдаёт:
Task reporting from thread 10, thread pool: False Task reporting from thread 11, thread pool: False Task reporting from thread 12, thread pool: False Task reporting from thread 13, thread pool: False Task reporting from thread 14, thread pool: False Task reporting from thread 15, thread pool: False Task reporting from thread 16, thread pool: False Task reporting from thread 9, thread pool: False Task reporting from thread 10, thread pool: False Task reporting from thread 11, thread pool: False