Цель: закрепить руками то, что разобрано в главах 1–17: писать обычный прикладной async-код (сбор данных по сети, обращения к базе, лента, очередь) и проверять себя автоматически.
Лабораторная:start/ — восемь методов-заготовок в Practice.cs и проверки; final/ — эталонные решения. Запуск: dotnet run -c Release --project start (все задачи) или -- 3 5 (только 3-ю и 5-ю).
Статус: ✅ эталон проходит все восемь проверок (три запуска подряд на Ubuntu 26.04, .NET 10.0.12 и три на Windows 11, 8 ядер: 8 из 8, времена проверок в тех же пределах). Проверки дополнительно прогнаны на заведомо неправильных решениях (последовательный цикл, неограниченный WhenAll, WhenAny + Delay и т. п.): каждая из восьми их отвергает с понятным сообщением.
Ни сети, ни базы данных не нужно. В Fakes.cs лежит «внешний мир» с настоящими задержками, сбоями и отменой:
Класс
Что изображает
Что умеет
FakeApi
HTTP-сервис
задержка LatencyMs; сбои FailFirstAttempts (первые N обращений к ключу падают с HttpRequestException), AlwaysFail, BadKeys; считает обращения, пик одновременных вызовов и отменённые вызовы
FakeDb
база данных
пакетная вставка InsertBatchAsync (30 мс); запоминает строки, размеры пакетов и пик одновременных записей
FakeFeed
постраничная лента
5 страниц по 10 элементов, 50 мс на страницу; считает запрошенные страницы
Правила:
Решайте только в start/Practice.cs. Fakes.cs и Checks.cs не меняйте.
Запустите dotnet run -c Release --project start. Пока метод не написан, проверка говорит «не реализовано».
Проверка пишет PASS/FAIL и причину провала словами («одновременно шло 40 запросов, лимит 5»).
Эталон смотрите после своей попытки. Он в final/Practice.cs и в раскрывающихся блоках ниже.
Задачи расположены от простого к сложному. Каждая опирается на главы курса; в каждой есть ловушка, которую проверка нарочно проверяет.
Проверка: десять запросов по 100 мс должны уложиться примерно в 100–200 мс (не в секунду); порядок результатов; ошибка любого запроса — исключение.
Подсказка
Сначала запустите все запросы (получите задачи), потом дождитесь их вместе. await внутри цикла запускает запросы по одному (глава 16, блок 6). Какой метод возвращает массив результатов в порядке аргументов?
Сценарий. Те же запросы, но их 40, а сервис разрешает не больше пяти одновременных соединений. Результаты по-прежнему в порядке ключей.
Сигнатура:Task<string[]> FetchLimitedAsync(FakeApi api, IReadOnlyList<string> keys, int maxParallel, CancellationToken ct).
Проверка: пик одновременных запросов ровно 5; общее время около 400 мс (8 «волн» по 50 мс); при отмене через 250 мс новые запросы не запускаются (из 100 ключей стартует меньше двадцати).
Подсказка
Есть три стандартных способа: Parallel.ForEachAsync с MaxDegreeOfParallelism, SemaphoreSlim (глава 10), несколько читателей одного Channel. Результаты складывайте по индексу, а не по порядку завершения.
Parallel.ForEachAsync сам держит не больше MaxDegreeOfParallelism задач и перестаёт брать новые после отмены.
Запись results[i] = ... в разные ячейки массива из разных потоков безопасна: каждый индекс пишет ровно один обработчик.
Вариант на SemaphoreSlim: await sem.WaitAsync(ct) перед запросом и Release() в finally, но тогда все 40 задач уже созданы и ждут у семафора; Parallel.ForEachAsync создаёт их по мере освобождения.
Требования: не уложились в таймаут — TimeoutException, и сам запрос отменён (а не брошен в фоне); отменил вызывающий — OperationCanceledException, а не TimeoutException.
Проверка: запрос на 2 с с таймаутом 200 мс падает с TimeoutException быстро, и у FakeApiCancelled == 1; быстрый запрос возвращает значение; внешняя отмена через 100 мс даёт OperationCanceledException.
Подсказка
Нужен один токен, который сработает и от внешней отмены, и от таймера (глава 9). Как отличить, кто отменил, когда пришёл OperationCanceledException?
usingvarcts=CancellationTokenSource.CreateLinkedTokenSource(ct);cts.CancelAfter(timeout);try{returnawaitapi.GetAsync(key,cts.Token);}catch(OperationCanceledException)when(!ct.IsCancellationRequested){thrownewTimeoutException($"{key}: не уложились в {timeout.TotalMilliseconds:F0} мс");}
CreateLinkedTokenSource + CancelAfter даёт токен «внешний ИЛИ таймер»; using освобождает связанный токен (глава 9: без Dispose обработчики копятся на родителе).
Фильтр when (!ct.IsCancellationRequested) различает причины: если внешний токен не отменён, отменил таймер.
Решение WhenAny(запрос, Delay(timeout)) здесь не пройдёт: оно возвращает управление по таймауту, но запрос продолжает работать (глава 9: «таймаут» без отмены). task.WaitAsync(timeout) тоже перестаёт ждать, но не отменяет запрос.
Требования: повтор только при HttpRequestException; всего не больше maxAttempts обращений; паузы между попытками baseDelay, 2×baseDelay, 4×baseDelay…; после последней неудачи пробрасывается её исключение; отмена прерывает и запрос, и паузу и не считается поводом для повтора.
Проверка: сервис падает дважды и отвечает на третий раз (3 обращения, суммарная пауза ≥ 300 мс); всегда падающий сервис — ровно 3 попытки и HttpRequestException; отмена через 150 мс во время паузы — OperationCanceledException без новых обращений.
Подсказка
Цикл с try/catch вокруг await. Условие when (attempt < maxAttempts) в catch избавляет от отдельной проверки «последняя ли попытка»: на последней исключение просто не перехватывается. Пауза — Task.Delayс токеном.
Ловим только HttpRequestException: OperationCanceledException не перехватывается и сразу выходит наружу.
Пауза Task.Delay(..., ct) отменяема. Без токена отмена во время паузы заметилась бы только после неё.
На последней попытке фильтр when ложный, исключение летит вызывающему как есть, со своим стеком.
В реальном коде добавляют случайный разброс (jitter) к паузе и ограничивают максимальную паузу; готовое решение — библиотека Polly или Microsoft.Extensions.Http.Resilience.
Сценарий. Один и тот же файл есть на трёх зеркалах. Спросите все сразу и возьмите первый успешный ответ; остальные запросы отмените, чтобы не тратить трафик.
Требования: быстро упавшее зеркало не должно «выиграть»; результат — первое успешное; остальные запросы отменены; если упали все — AggregateException со всеми ошибками.
Проверка: зеркало 0 падает за 50 мс, зеркало 1 отвечает за 150 мс, зеркало 2 — за 2 с. Ждём ответ зеркала 1 за 150–500 мс и Cancelled == 1 у зеркала 2. Если падают все три, AggregateException с тремя ошибками.
Подсказка
Task.WhenAny вернёт первую завершённую задачу, а она может быть упавшей. Нужен цикл: ждём очередную, смотрим на результат, упавшую убираем из списка, пока не кончатся. Отмена остальных — через общий связанный токен.
Цикл WhenAny + Remove для небольшого числа задач — нормально. Для тысяч задач он квадратичен (глава 10: 16 000 задач заняли около 2 с); тогда Task.WhenEach (.NET 9+).
cts.Cancel() отменяет все оставшиеся запросы; они завершатся отменой, но это нас больше не касается: возвращённый результат уже взят.
Исключения отменённых задач никто не наблюдает, и это безопасно: отменённая задача не попадает в UnobservedTaskException.
Строка ct.ThrowIfCancellationRequested() перед throw: если упали все из-за внешней отмены, правильный ответ — отмена, а не AggregateException.
Сценарий. Сотня одновременных запросов просит один и тот же справочник. Загружать его надо один раз, остальные ждут тот же результат. Если загрузка упала, ошибку нельзя запоминать навсегда.
Класс:Practice.AsyncCache(Func<string, Task<string>> factory) с методом Task<string> GetAsync(string key).
Проверка: 20 одновременных GetAsync("a") вызывают factory один раз; готовое значение не пересчитывается; другой ключ загружается отдельно; если первая загрузка упала, следующий вызов пробует заново.
Подсказка
ConcurrentDictionary.GetOrAdd(key, factory) может вызвать factoryнесколько раз при гонке (выигрывает одна запись, но остальные вызовы уже сделаны). Что можно положить в словарь вместо готового значения, чтобы запуск загрузки произошёл ровно один раз? Вспомните Lazy<T> и AsyncLazy из главы 10.
privatereadonlyConcurrentDictionary<string,Lazy<Task<string>>>_items=new();publicTask<string>GetAsync(stringkey){varlazy=_items.GetOrAdd(key,k=>{Lazy<Task<string>>?self=null;self=newLazy<Task<string>>(()=>LoadAsync(k,self!));returnself;});returnlazy.Value;}privateasyncTask<string>LoadAsync(stringkey,Lazy<Task<string>>self){try{returnawaitfactory(key);}catch{_items.TryRemove(newKeyValuePair<string,Lazy<Task<string>>>(key,self));// только свою записьthrow;}}
В словаре лежит Lazy<Task<string>>: GetOrAdd может создать несколько Lazy, но .Value у проигравших не вызывается, поэтому фабрика запускается один раз.
Все ожидающие получают один и тот же Task, то есть одно значение или одну ошибку.
При ошибке мы удаляем именно свою запись (TryRemove(KeyValuePair)), чтобы не стереть более новую запись, появившуюся после сбоя.
Lazy<Task<T>> без удаления при ошибке кеширует сбой навсегда (глава 10, «Lazy<Task> кеширует ошибку»).
Чего здесь нет: отмены. Если один из ожидающих отменится, общую загрузку отменять нельзя, потому что она нужна остальным; поэтому фабрика токена не получает.
Сценарий. Лента отдаётся страницами по 10 элементов. Нужен IAsyncEnumerable<string>, из которого потребитель берёт элементы по одному и может остановиться в любой момент.
Требования: страницы читаются по мере надобности (следующую не запрашивать, пока потребитель не дошёл до конца текущей); отмена через WithCancellation(token) прерывает чтение.
Проверка: полный проход — 50 элементов по порядку и ровно 5 запросов страниц; потребитель остановился на 15-м элементе — запрошено не больше 2 страниц; WithCancellation с токеном, отменяемым через 80 мс, даёт OperationCanceledException.
Подсказка
Асинхронный итератор: async IAsyncEnumerable<T> с yield return (глава 11). Чтобы токен из WithCancellation попал в ваш метод, параметр нужно пометить атрибутом [EnumeratorCancellation].
Итератор ленив: код после yield return выполняется только когда потребитель запрашивает следующий элемент, поэтому «читать вперёд» не получается само собой.
Без [EnumeratorCancellation] токен, переданный через WithCancellation, до вашего ct не дойдёт, и отмена будет работать только между элементами.
HasNext избавляет от лишнего запроса пустой страницы в конце.
Если нужно предзагружать следующую страницу в фоне (быстрее, но тратит запросы), это уже Channel с ёмкостью 1: см. задачу 8.
Сценарий. Нужно обработать 95 ключей: получить данные из API (не больше 4 одновременных запросов) и записать в базу пакетами по 10. Запись в базу — одним писателем (так проще и безопаснее для транзакций). Очереди между стадиями ограничены, чтобы медленная база не раздувала память.
Сигнатура:Task<int> RunPipelineAsync(FakeApi api, FakeDb db, IEnumerable<string> keys, int workers, int batchSize, CancellationToken ct) — возвращает число записанных строк.
Требования: все строки записаны ровно по разу; пакеты по batchSize, последний может быть меньше; не больше workers запросов одновременно; ошибка любой стадии останавливает остальные, а наружу выходит настоящая причина (а не OperationCanceledException от остановленной стадии).
Проверка: 95 ключей → 95 строк, девять пакетов по 10 и один из 5, пик запросов ≤ 4, пик записей = 1; если API падает на ключе k30, метод бросает HttpRequestException, быстро и до обработки всех ключей.
Подсказка
Три стадии, две ограниченные очереди: производитель → input → workers читателей → fetched → один писатель. Каждая стадия — отдельная задача. Все используют один CancellationTokenSource; любая стадия при ошибке отменяет его. Не забудьте закрыть (Complete) fetched, когда все читатели закончили.
usingvarcts=CancellationTokenSource.CreateLinkedTokenSource(ct);varinput=Channel.CreateBounded<string>(20);varfetched=Channel.CreateBounded<string>(20);asyncTaskGuard(Func<Task>stage){try{awaitstage();}catch{cts.Cancel();throw;}// ошибка одной стадии останавливает остальные}varproducer=Guard(async()=>{foreach(varkeyinkeys)awaitinput.Writer.WriteAsync(key,cts.Token);input.Writer.Complete();});varfetchers=Enumerable.Range(0,workers).Select(_=>Guard(async()=>{awaitforeach(varkeyininput.Reader.ReadAllAsync(cts.Token))awaitfetched.Writer.WriteAsync(awaitapi.GetAsync(key,cts.Token),cts.Token);})).ToArray();varcloser=Guard(async()=>{awaitTask.WhenAll(fetchers);fetched.Writer.Complete();});inttotal=0;varwriter=Guard(async()=>{varbatch=newList<string>(batchSize);awaitforeach(varrowinfetched.Reader.ReadAllAsync(cts.Token)){batch.Add(row);if(batch.Count<batchSize)continue;awaitdb.InsertBatchAsync(batch.ToArray(),cts.Token);total+=batch.Count;batch.Clear();}if(batch.Count>0){awaitdb.InsertBatchAsync(batch.ToArray(),cts.Token);total+=batch.Count;}});varall=Task.WhenAll(producer,closer,writer);try{awaitall;}catch{// WhenAll бросит первое по порядку исключение, а оно может быть OperationCanceledException от «остановленной» стадии.varreal=all.Exception!.InnerExceptions.FirstOrDefault(e=>eisnotOperationCanceledException)??all.Exception.InnerExceptions[0];ExceptionDispatchInfo.Capture(real).Throw();}returntotal;
Channel.CreateBounded(20) даёт обратное давление: если база медленная, fetched заполняется, читатели ждут в WriteAsync, затем заполняется input, и производитель замедляется (глава 11).
closer завершает fetched, только когда закончили все читатели (WhenAll). Закрыть канал в каждом читателе нельзя: первый же завершившийся оборвал бы остальных.
Guard отменяет общий токен при любой ошибке, и остальные стадии выходят из WriteAsync/ReadAllAsync по отмене.
Блок catch после WhenAll решает известную ловушку: await бросает первое по порядку исключение, а им может оказаться OperationCanceledException от стадии, которую остановили из-за настоящей ошибки. Поэтому берём из all.Exception первое исключение, не являющееся отменой (глава 8).
Один писатель в базу гарантирует и порядок пакетов, и то, что db.PeakConcurrentWrites == 1.
Вариант попроще: Parallel.ForEachAsync для запросов и пакеты в потокобезопасный буфер. Он закрывает ограничение запросов, но обратное давление от медленной базы и единственного писателя придётся добавлять отдельно.
usingSystem.Runtime.CompilerServices;// Восемь задач. Тела методов пустые (NotImplementedException): напишите их сами. Эталонные решения — в final/. Описание задач и разбор — в README главы.staticclassPractice{/// <summary>/// Задача 1. Получить данные по всем ключам <b>параллельно</b>. Порядок результатов — как порядок ключей./// Если любой запрос упал — исключение этого запроса./// </summary>publicstaticTask<string[]>FetchAllAsync(FakeApiapi,IReadOnlyList<string>keys,CancellationTokenct){// TODO: ваше решениеthrownewNotImplementedException();}/// <summary>/// Задача 2. То же, но одновременно выполняется не больше <paramref name="maxParallel"/> запросов./// Порядок результатов — как порядок ключей. При отмене новые запросы не запускаются./// </summary>publicstaticasyncTask<string[]>FetchLimitedAsync(FakeApiapi,IReadOnlyList<string>keys,intmaxParallel,CancellationTokenct){// TODO: ваше решениеthrownewNotImplementedException();}/// <summary>/// Задача 3. Запрос с таймаутом. Не уложились — <see cref="TimeoutException"/>, и сам запрос должен быть отменён./// Если отменил вызывающий (<paramref name="ct"/>) — <see cref="OperationCanceledException"/>, а не TimeoutException./// </summary>publicstaticasyncTask<string>GetWithTimeoutAsync(FakeApiapi,stringkey,TimeSpantimeout,CancellationTokenct){// TODO: ваше решениеthrownewNotImplementedException();}/// <summary>/// Задача 4. Повтор при <see cref="HttpRequestException"/>: до <paramref name="maxAttempts"/> обращений,/// между попытками пауза baseDelay, 2×baseDelay, 4×baseDelay… Последняя ошибка пробрасывается./// Отмена прерывает и запрос, и паузу и повторов не вызывает./// </summary>publicstaticasyncTask<string>GetWithRetryAsync(FakeApiapi,stringkey,intmaxAttempts,TimeSpanbaseDelay,CancellationTokenct){// TODO: ваше решениеthrownewNotImplementedException();}/// <summary>/// Задача 5. Одновременно спросить все зеркала и вернуть ответ первого <b>успешного</b>. Остальные запросы отменить./// Если упали все — <see cref="AggregateException"/> со всеми ошибками./// </summary>publicstaticasyncTask<string>FirstSuccessAsync(IReadOnlyList<FakeApi>mirrors,stringkey,CancellationTokenct){// TODO: ваше решениеthrownewNotImplementedException();}/// <summary>/// Задача 6. Асинхронный кеш. Одновременные запросы одного ключа вызывают фабрику <b>один раз</b>;/// готовое значение хранится. Упавшая загрузка в кеше не остаётся: следующий запрос пробует заново./// </summary>publicsealedclassAsyncCache(Func<string,Task<string>>factory){publicTask<string>GetAsync(stringkey){// TODO: ваше решение (factory(key) вызывается в вашем коде)thrownewNotImplementedException();}}/// <summary>/// Задача 7. Лента как <see cref="IAsyncEnumerable{T}"/>: страницы читаются по мере надобности (страницу N+1 не запрашивать,/// пока потребитель не дошёл до конца страницы N). Отмена через WithCancellation должна прерывать чтение./// </summary>publicstaticasyncIAsyncEnumerable<string>ReadFeedAsync(FakeFeedfeed,[EnumeratorCancellation]CancellationTokenct=default){// TODO: ваше решение (в теле метода должен быть yield return)thrownewNotImplementedException();#pragma warning disable CS0162yieldbreak;// заглушка, чтобы метод компилировался; удалите вместе с pragma#pragma warning restore CS0162}/// <summary>/// Задача 8. Конвейер «ключи → API → база»: <paramref name="workers"/> параллельных запросов к API,/// запись в базу одним писателем пакетами по <paramref name="batchSize"/> (последний пакет может быть меньше)./// Очереди между стадиями ограничены (обратное давление). Ошибка любой стадии отменяет остальные и/// пробрасывается наружу (настоящая причина, а не OperationCanceledException). Возвращает число записанных строк./// </summary>publicstaticasyncTask<int>RunPipelineAsync(FakeApiapi,FakeDbdb,IEnumerable<string>keys,intworkers,intbatchSize,CancellationTokenct){// TODO: ваше решениеthrownewNotImplementedException();}}
usingSystem.Collections.Concurrent;usingSystem.Runtime.CompilerServices;usingSystem.Runtime.ExceptionServices;usingSystem.Threading.Channels;// Восемь задач. В start/ тела методов пустые (NotImplementedException): напишите их сами.// В final/ лежат эталонные решения. Описание задач и разбор — в README главы.staticclassPractice{/// <summary>/// Задача 1. Получить данные по всем ключам <b>параллельно</b>. Порядок результатов — как порядок ключей./// Если любой запрос упал — исключение этого запроса./// </summary>publicstaticTask<string[]>FetchAllAsync(FakeApiapi,IReadOnlyList<string>keys,CancellationTokenct){// <решение>returnTask.WhenAll(keys.Select(k=>api.GetAsync(k,ct)));// </решение>}/// <summary>/// Задача 2. То же, но одновременно выполняется не больше <paramref name="maxParallel"/> запросов./// Порядок результатов — как порядок ключей. При отмене новые запросы не запускаются./// </summary>publicstaticasyncTask<string[]>FetchLimitedAsync(FakeApiapi,IReadOnlyList<string>keys,intmaxParallel,CancellationTokenct){// <решение>varresults=newstring[keys.Count];awaitParallel.ForEachAsync(Enumerable.Range(0,keys.Count),newParallelOptions{MaxDegreeOfParallelism=maxParallel,CancellationToken=ct},async(i,token)=>results[i]=awaitapi.GetAsync(keys[i],token));returnresults;// </решение>}/// <summary>/// Задача 3. Запрос с таймаутом. Не уложились — <see cref="TimeoutException"/>, и сам запрос должен быть отменён./// Если отменил вызывающий (<paramref name="ct"/>) — <see cref="OperationCanceledException"/>, а не TimeoutException./// </summary>publicstaticasyncTask<string>GetWithTimeoutAsync(FakeApiapi,stringkey,TimeSpantimeout,CancellationTokenct){// <решение>usingvarcts=CancellationTokenSource.CreateLinkedTokenSource(ct);cts.CancelAfter(timeout);try{returnawaitapi.GetAsync(key,cts.Token);}catch(OperationCanceledException)when(!ct.IsCancellationRequested){thrownewTimeoutException($"{key}: не уложились в {timeout.TotalMilliseconds:F0} мс");}// </решение>}/// <summary>/// Задача 4. Повтор при <see cref="HttpRequestException"/>: до <paramref name="maxAttempts"/> обращений,/// между попытками пауза baseDelay, 2×baseDelay, 4×baseDelay… Последняя ошибка пробрасывается./// Отмена прерывает и запрос, и паузу и повторов не вызывает./// </summary>publicstaticasyncTask<string>GetWithRetryAsync(FakeApiapi,stringkey,intmaxAttempts,TimeSpanbaseDelay,CancellationTokenct){// <решение>for(intattempt=1;;attempt++){try{returnawaitapi.GetAsync(key,ct);}catch(HttpRequestException)when(attempt<maxAttempts){awaitTask.Delay(baseDelay*(1<<(attempt-1)),ct);}}// </решение>}/// <summary>/// Задача 5. Одновременно спросить все зеркала и вернуть ответ первого <b>успешного</b>. Остальные запросы отменить./// Если упали все — <see cref="AggregateException"/> со всеми ошибками./// </summary>publicstaticasyncTask<string>FirstSuccessAsync(IReadOnlyList<FakeApi>mirrors,stringkey,CancellationTokenct){// <решение>usingvarcts=CancellationTokenSource.CreateLinkedTokenSource(ct);varpending=mirrors.Select(m=>m.GetAsync(key,cts.Token)).ToList();varerrors=newList<Exception>();while(pending.Count>0){Task<string>done=awaitTask.WhenAny(pending);pending.Remove(done);if(done.IsCompletedSuccessfully){cts.Cancel();returndone.Result;}errors.Add(done.Exception!.InnerException??done.Exception);}ct.ThrowIfCancellationRequested();thrownewAggregateException(errors);// </решение>}/// <summary>/// Задача 6. Асинхронный кеш. Одновременные запросы одного ключа вызывают фабрику <b>один раз</b>;/// готовое значение хранится. Упавшая загрузка в кеше не остаётся: следующий запрос пробует заново./// </summary>publicsealedclassAsyncCache(Func<string,Task<string>>factory){// <решение>privatereadonlyConcurrentDictionary<string,Lazy<Task<string>>>_items=new();publicTask<string>GetAsync(stringkey){varlazy=_items.GetOrAdd(key,k=>{Lazy<Task<string>>?self=null;self=newLazy<Task<string>>(()=>LoadAsync(k,self!));returnself;});returnlazy.Value;}privateasyncTask<string>LoadAsync(stringkey,Lazy<Task<string>>self){try{returnawaitfactory(key);}catch{_items.TryRemove(newKeyValuePair<string,Lazy<Task<string>>>(key,self));// только свою записьthrow;}}// </решение>}/// <summary>/// Задача 7. Лента как <see cref="IAsyncEnumerable{T}"/>: страницы читаются по мере надобности (страницу N+1 не запрашивать,/// пока потребитель не дошёл до конца страницы N). Отмена через WithCancellation должна прерывать чтение./// </summary>publicstaticasyncIAsyncEnumerable<string>ReadFeedAsync(FakeFeedfeed,[EnumeratorCancellation]CancellationTokenct=default){// <решение>for(intpage=0;;page++){Pagep=awaitfeed.GetPageAsync(page,ct);foreach(variteminp.Items)yieldreturnitem;if(!p.HasNext)yieldbreak;}// </решение>}/// <summary>/// Задача 8. Конвейер «ключи → API → база»: <paramref name="workers"/> параллельных запросов к API,/// запись в базу одним писателем пакетами по <paramref name="batchSize"/> (последний пакет может быть меньше)./// Очереди между стадиями ограничены (обратное давление). Ошибка любой стадии отменяет остальные и/// пробрасывается наружу (настоящая причина, а не OperationCanceledException). Возвращает число записанных строк./// </summary>publicstaticasyncTask<int>RunPipelineAsync(FakeApiapi,FakeDbdb,IEnumerable<string>keys,intworkers,intbatchSize,CancellationTokenct){// <решение>usingvarcts=CancellationTokenSource.CreateLinkedTokenSource(ct);varinput=Channel.CreateBounded<string>(20);varfetched=Channel.CreateBounded<string>(20);asyncTaskGuard(Func<Task>stage){try{awaitstage();}catch{cts.Cancel();throw;}// ошибка одной стадии останавливает остальные}varproducer=Guard(async()=>{foreach(varkeyinkeys)awaitinput.Writer.WriteAsync(key,cts.Token);input.Writer.Complete();});varfetchers=Enumerable.Range(0,workers).Select(_=>Guard(async()=>{awaitforeach(varkeyininput.Reader.ReadAllAsync(cts.Token))awaitfetched.Writer.WriteAsync(awaitapi.GetAsync(key,cts.Token),cts.Token);})).ToArray();varcloser=Guard(async()=>{awaitTask.WhenAll(fetchers);fetched.Writer.Complete();});inttotal=0;varwriter=Guard(async()=>{varbatch=newList<string>(batchSize);awaitforeach(varrowinfetched.Reader.ReadAllAsync(cts.Token)){batch.Add(row);if(batch.Count<batchSize)continue;awaitdb.InsertBatchAsync(batch.ToArray(),cts.Token);total+=batch.Count;batch.Clear();}if(batch.Count>0){awaitdb.InsertBatchAsync(batch.ToArray(),cts.Token);total+=batch.Count;}});varall=Task.WhenAll(producer,closer,writer);try{awaitall;}catch{// WhenAll бросит первое по порядку исключение, а оно может быть OperationCanceledException от «остановленной» стадии.varreal=all.Exception!.InnerExceptions.FirstOrDefault(e=>eisnotOperationCanceledException)??all.Exception.InnerExceptions[0];ExceptionDispatchInfo.Capture(real).Throw();}returntotal;// </решение>}}
usingSystem.Collections.Concurrent;// «Внешний мир» для задач: ни сети, ни базы не нужно, но задержки, сбои и отмена ведут себя как настоящие.// Эти файлы трогать не надо: решайте в Practice.cs./// <summary>Условный HTTP-сервис. Считает вызовы, одновременные вызовы и отменённые вызовы.</summary>sealedclassFakeApi(intlatencyMs=100){privateint_inFlight,_peak,_calls,_cancelled;privatereadonlyConcurrentDictionary<string,int>_attempts=new();publicintLatencyMs{get;set;}=latencyMs;/// <summary>Первые N обращений к каждому ключу падают с HttpRequestException (503).</summary>publicintFailFirstAttempts{get;set;}publicboolAlwaysFail{get;set;}publicHashSet<string>BadKeys{get;}=[];publicintCalls=>Volatile.Read(ref_calls);publicintPeak=>Volatile.Read(ref_peak);publicintCancelled=>Volatile.Read(ref_cancelled);publicasyncTask<string>GetAsync(stringkey,CancellationTokenct=default){Interlocked.Increment(ref_calls);intnow=Interlocked.Increment(ref_inFlight);intseen;while(now>(seen=Volatile.Read(ref_peak))&&Interlocked.CompareExchange(ref_peak,now,seen)!=seen){}try{awaitTask.Delay(LatencyMs,ct);intattempt=_attempts.AddOrUpdate(key,1,(_,n)=>n+1);if(AlwaysFail||BadKeys.Contains(key)||attempt<=FailFirstAttempts)thrownewHttpRequestException($"503 для {key} (попытка {attempt})");return$"data:{key}";}catch(OperationCanceledException){Interlocked.Increment(ref_cancelled);throw;}finally{Interlocked.Decrement(ref_inFlight);}}}/// <summary>Условная база данных: пакетная вставка. Запоминает всё записанное.</summary>sealedclassFakeDb{privateint_writing,_peakWriting;publicConcurrentBag<string>Rows{get;}=[];publicConcurrentQueue<int>BatchSizes{get;}=[];publicintPeakConcurrentWrites=>Volatile.Read(ref_peakWriting);publicasyncTaskInsertBatchAsync(IReadOnlyCollection<string>rows,CancellationTokenct=default){intnow=Interlocked.Increment(ref_writing);intseen;while(now>(seen=Volatile.Read(ref_peakWriting))&&Interlocked.CompareExchange(ref_peakWriting,now,seen)!=seen){}try{awaitTask.Delay(30,ct);BatchSizes.Enqueue(rows.Count);foreach(varrinrows)Rows.Add(r);}finally{Interlocked.Decrement(ref_writing);}}}sealedrecordPage(IReadOnlyList<string>Items,boolHasNext);/// <summary>Постраничная лента: 5 страниц по 10 элементов.</summary>sealedclassFakeFeed{privateint_pages;publicintPagesRequested=>Volatile.Read(ref_pages);publicasyncTask<Page>GetPageAsync(intpage,CancellationTokenct=default){Interlocked.Increment(ref_pages);awaitTask.Delay(50,ct);varitems=Enumerable.Range(page*10,10).Select(i=>$"item{i}").ToList();returnnewPage(items,page<4);}}
usingSystem.Diagnostics;// Проверки для самопроверки. Каждая запускает ваше решение из Practice.cs на «внешнем мире» из Fakes.cs.staticclassChecks{publicstaticasyncTaskRunAsync(string[]args){varall=new(intn,stringtitle,Func<Task<string?>>run)[]{(1,"Параллельный сбор",Check1),(2,"Ограничение параллелизма",Check2),(3,"Таймаут и отмена",Check3),(4,"Повтор с экспоненциальной паузой",Check4),(5,"Первый успешный из зеркал",Check5),(6,"Кеш без дублей",Check6),(7,"Постраничная лента как поток",Check7),(8,"Конвейер на Channel",Check8),};intpassed=0,ran=0;foreach(var(n,title,run)inall.Where(t=>args.Length==0||args.Contains(t.n.ToString()))){ran++;varsw=Stopwatch.StartNew();string?error;try{vartask=run();error=awaitTask.WhenAny(task,Task.Delay(15_000))==task?awaittask:"не уложились в 15 секунд (зависло?)";}catch(NotImplementedException){error="не реализовано";}catch(Exceptione){error=$"исключение {e.GetType().Name}: {e.Message}";}if(errorisnull)passed++;Console.WriteLine($"{(error is null ? "PASS" : "FAIL")} {n}. {title,-34} {sw.ElapsedMilliseconds,5} мс{(error is null ? "" : "→" + error)}");}Console.WriteLine($"\nИтого: {passed} из {ran}");}staticstring[]Keys(intn)=>Enumerable.Range(0,n).Select(i=>$"k{i}").ToArray();staticstring[]Expected(IEnumerable<string>keys)=>keys.Select(k=>$"data:{k}").ToArray();// ---- 1 ----staticasyncTask<string?>Check1(){varapi=newFakeApi(100);varkeys=Keys(10);varsw=Stopwatch.StartNew();varr=awaitPractice.FetchAllAsync(api,keys,CancellationToken.None);if(!r.SequenceEqual(Expected(keys)))return"результаты не совпали или порядок нарушен";if(sw.ElapsedMilliseconds>450)return$"10 запросов по 100 мс заняли {sw.ElapsedMilliseconds} мс: они идут не параллельно";returnnull;}// ---- 2 ----staticasyncTask<string?>Check2(){varapi=newFakeApi(50);varkeys=Keys(40);varsw=Stopwatch.StartNew();varr=awaitPractice.FetchLimitedAsync(api,keys,5,CancellationToken.None);if(!r.SequenceEqual(Expected(keys)))return"результаты не совпали или порядок нарушен";if(api.Peak>5)return$"одновременно шло {api.Peak} запросов, лимит 5";if(api.Peak<5)return$"одновременно шло только {api.Peak}: лимит 5 не использован";if(sw.ElapsedMilliseconds>1200)return$"слишком медленно: {sw.ElapsedMilliseconds} мс (ожидали около 400)";varapi2=newFakeApi(100);usingvarcts=newCancellationTokenSource(250);try{awaitPractice.FetchLimitedAsync(api2,Keys(100),5,cts.Token);return"при отмене не было исключения";}catch(OperationCanceledException){}if(api2.Calls>20)return$"после отмены продолжали запускаться запросы: всего {api2.Calls}";returnnull;}// ---- 3 ----staticasyncTask<string?>Check3(){varslow=newFakeApi(2000);varsw=Stopwatch.StartNew();try{awaitPractice.GetWithTimeoutAsync(slow,"a",TimeSpan.FromMilliseconds(200),CancellationToken.None);return"таймаут не сработал";}catch(TimeoutException){}if(sw.ElapsedMilliseconds>700)return$"таймаут 200 мс сработал через {sw.ElapsedMilliseconds} мс";awaitTask.Delay(100);if(slow.Cancelled!=1)return$"сам запрос не отменён (отменено {slow.Cancelled}): после таймаута он продолжает работать впустую";varfast=newFakeApi(50);if(awaitPractice.GetWithTimeoutAsync(fast,"b",TimeSpan.FromSeconds(1),CancellationToken.None)!="data:b")return"быстрый запрос вернул не то";usingvarcts=newCancellationTokenSource(100);try{awaitPractice.GetWithTimeoutAsync(newFakeApi(2000),"c",TimeSpan.FromSeconds(5),cts.Token);return"внешняя отмена не сработала";}catch(TimeoutException){return"внешняя отмена превратилась в TimeoutException: различайте причины";}catch(OperationCanceledException){}returnnull;}// ---- 4 ----staticasyncTask<string?>Check4(){varapi=newFakeApi(10){FailFirstAttempts=2};varsw=Stopwatch.StartNew();varr=awaitPractice.GetWithRetryAsync(api,"x",4,TimeSpan.FromMilliseconds(100),CancellationToken.None);if(r!="data:x")return"результат не тот";if(api.Calls!=3)return$"ожидали 3 обращения (2 сбоя и успех), было {api.Calls}";if(sw.ElapsedMilliseconds<290)return$"паузы слишком короткие: {sw.ElapsedMilliseconds} мс (нужно 100 + 200 мс между попытками)";if(sw.ElapsedMilliseconds>900)return$"слишком долго: {sw.ElapsedMilliseconds} мс";vardead=newFakeApi(10){AlwaysFail=true};try{awaitPractice.GetWithRetryAsync(dead,"y",3,TimeSpan.FromMilliseconds(20),CancellationToken.None);return"ошибка проглочена";}catch(HttpRequestException){}if(dead.Calls!=3)return$"ожидали ровно 3 попытки, было {dead.Calls}";varapi3=newFakeApi(10){AlwaysFail=true};usingvarcts=newCancellationTokenSource(150);try{awaitPractice.GetWithRetryAsync(api3,"z",10,TimeSpan.FromMilliseconds(100),cts.Token);return"при отмене не было исключения";}catch(OperationCanceledException){}catch(HttpRequestException){return"отмена во время паузы превратилась в HttpRequestException или повторы не прерываются";}if(api3.Calls>3)return$"после отмены повторы продолжались: {api3.Calls} обращений";returnnull;}// ---- 5 ----staticasyncTask<string?>Check5(){varm0=newFakeApi(50){AlwaysFail=true};varm1=newFakeApi(150);varm2=newFakeApi(2000);varsw=Stopwatch.StartNew();varr=awaitPractice.FirstSuccessAsync([m0,m1,m2],"k",CancellationToken.None);if(r!="data:k")return"результат не тот";if(sw.ElapsedMilliseconds>500)return$"ждали {sw.ElapsedMilliseconds} мс: должен вернуться первый успешный (около 150 мс)";awaitTask.Delay(100);if(m2.Cancelled!=1)return"медленное зеркало не отменено: запрос продолжает работать впустую";try{awaitPractice.FirstSuccessAsync([newFakeApi(20){AlwaysFail=true},newFakeApi(30){AlwaysFail=true},newFakeApi(40){AlwaysFail=true}],"k",CancellationToken.None);return"когда падают все, нужно исключение";}catch(AggregateExceptione)when(e.InnerExceptions.Count==3){}catch(AggregateExceptione){return$"AggregateException должно содержать все 3 ошибки, а в нём {e.InnerExceptions.Count}";}returnnull;}// ---- 6 ----staticasyncTask<string?>Check6(){intcalls=0;varcache=newPractice.AsyncCache(asynckey=>{Interlocked.Increment(refcalls);awaitTask.Delay(100);return$"v:{key}";});varresults=awaitTask.WhenAll(Enumerable.Range(0,20).Select(_=>cache.GetAsync("a")));if(results.Any(x=>x!="v:a"))return"значения не совпали";if(calls!=1)return$"20 одновременных запросов одного ключа вызвали фабрику {calls} раз (нужно 1)";awaitcache.GetAsync("a");if(calls!=1)return"повторный запрос готового ключа снова вызвал фабрику";awaitcache.GetAsync("b");if(calls!=2)return"другой ключ должен вызывать фабрику отдельно";intattempts=0;varflaky=newPractice.AsyncCache(asynckey=>{intn=Interlocked.Increment(refattempts);awaitTask.Delay(50);if(n==1)thrownewInvalidOperationException("первый вызов падает");return"ok";});try{awaitflaky.GetAsync("x");return"первый вызов должен был упасть";}catch(InvalidOperationException){}try{if(awaitflaky.GetAsync("x")!="ok")return"после сбоя вернулось не то значение";}catch(InvalidOperationException){return"упавшая загрузка осталась в кеше: следующий вызов получил ту же ошибку (нужно пробовать заново)";}returnnull;}// ---- 7 ----staticasyncTask<string?>Check7(){varfeed=newFakeFeed();varall=newList<string>();awaitforeach(varxinPractice.ReadFeedAsync(feed))all.Add(x);if(all.Count!=50||!all.SequenceEqual(Enumerable.Range(0,50).Select(i=>$"item{i}")))return$"ожидали 50 элементов по порядку, получено {all.Count}";if(feed.PagesRequested!=5)return$"ожидали 5 запросов страниц, было {feed.PagesRequested}";varfeed2=newFakeFeed();varfirstFifteen=newList<string>();awaitforeach(varxinPractice.ReadFeedAsync(feed2)){firstFifteen.Add(x);if(firstFifteen.Count==15)break;}awaitTask.Delay(150);if(feed2.PagesRequested>2)return$"потребитель остановился на 15-м элементе, а страниц запрошено {feed2.PagesRequested} (нужно 2): лента читает вперёд";varfeed3=newFakeFeed();usingvarcts=newCancellationTokenSource(80);try{awaitforeach(var_inPractice.ReadFeedAsync(feed3).WithCancellation(cts.Token)){}return"WithCancellation(token) не прервал чтение: нужен [EnumeratorCancellation]";}catch(OperationCanceledException){}returnnull;}// ---- 8 ----staticasyncTask<string?>Check8(){varapi=newFakeApi(20);vardb=newFakeDb();varkeys=Keys(95);varsw=Stopwatch.StartNew();intn=awaitPractice.RunPipelineAsync(api,db,keys,workers:4,batchSize:10,CancellationToken.None);if(n!=95)return$"вернули {n}, ожидали 95";if(!db.Rows.OrderBy(x=>x).SequenceEqual(Expected(keys).OrderBy(x=>x)))return"в базе не все записи или есть дубли";if(api.Peak>4)return$"одновременно {api.Peak} запросов к API, воркеров 4";if(db.PeakConcurrentWrites!=1)return$"запись в базу шла в {db.PeakConcurrentWrites} потока(ов) одновременно; ожидали одного писателя";varsizes=db.BatchSizes.ToArray();if(sizes.Count(s=>s==10)!=9||sizes.Count(s=>s==5)!=1||sizes.Length!=10)return$"пакеты: {string.Join(",", sizes)}; ожидали девять по 10 и один из 5";if(sw.ElapsedMilliseconds>3000)return$"слишком медленно: {sw.ElapsedMilliseconds} мс";varapi2=newFakeApi(20);api2.BadKeys.Add("k30");vardb2=newFakeDb();varsw2=Stopwatch.StartNew();try{awaitPractice.RunPipelineAsync(api2,db2,Keys(95),4,10,CancellationToken.None);return"ошибка на одном ключе должна ломать конвейер";}catch(HttpRequestException){}catch(Exceptione){return$"ожидали HttpRequestException (настоящую причину), получили {e.GetType().Name}";}if(sw2.ElapsedMilliseconds>3000)return"после ошибки конвейер завершался слишком долго";if(api2.Calls>=95)return"после ошибки конвейер продолжал обрабатывать все ключи: остановите остальных";returnnull;}}
// Глава 20. Практикум: восемь задач на async/await.// Решайте в Practice.cs, проверяйте:// dotnet run -c Release --project start — все задачи// dotnet run -c Release --project start -- 3 5 — только задачи 3 и 5awaitChecks.RunAsync(args);