Перейти к содержанию

20. Практикум: восемь задач на async/await

О главе

Цель: закрепить руками то, что разобрано в главах 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 и т. п.): каждая из восьми их отвергает с понятным сообщением.

20.1. Как работать

Ни сети, ни базы данных не нужно. В Fakes.cs лежит «внешний мир» с настоящими задержками, сбоями и отменой:

Класс Что изображает Что умеет
FakeApi HTTP-сервис задержка LatencyMs; сбои FailFirstAttempts (первые N обращений к ключу падают с HttpRequestException), AlwaysFail, BadKeys; считает обращения, пик одновременных вызовов и отменённые вызовы
FakeDb база данных пакетная вставка InsertBatchAsync (30 мс); запоминает строки, размеры пакетов и пик одновременных записей
FakeFeed постраничная лента 5 страниц по 10 элементов, 50 мс на страницу; считает запрошенные страницы

Правила:

  1. Решайте только в start/Practice.cs. Fakes.cs и Checks.cs не меняйте.
  2. Запустите dotnet run -c Release --project start. Пока метод не написан, проверка говорит «не реализовано».
  3. Проверка пишет PASS/FAIL и причину провала словами («одновременно шло 40 запросов, лимит 5»).
  4. Эталон смотрите после своей попытки. Он в final/Practice.cs и в раскрывающихся блоках ниже.

Задачи расположены от простого к сложному. Каждая опирается на главы курса; в каждой есть ловушка, которую проверка нарочно проверяет.

20.2. Задачи

Задача 1. Параллельный сбор

Сценарий. Страница профиля собирается из десяти независимых запросов. Нужно получить все и вернуть результаты в порядке ключей.

Сигнатура: Task<string[]> FetchAllAsync(FakeApi api, IReadOnlyList<string> keys, CancellationToken ct).

Проверка: десять запросов по 100 мс должны уложиться примерно в 100–200 мс (не в секунду); порядок результатов; ошибка любого запроса — исключение.

Подсказка

Сначала запустите все запросы (получите задачи), потом дождитесь их вместе. await внутри цикла запускает запросы по одному (глава 16, блок 6). Какой метод возвращает массив результатов в порядке аргументов?

Эталон и разбор
final/Practice.cs
        return Task.WhenAll(keys.Select(k => api.GetAsync(k, ct)));
  • Select лениво создаёт задачи, а Task.WhenAll перечисляет их все сразу: к моменту ожидания все запросы уже идут (в отличие от foreach с await).
  • WhenAll возвращает результаты в порядке входа, а не в порядке завершения.
  • Если упадёт несколько запросов, await бросит только первое исключение, остальные лежат в Task.Exception (глава 8).

Задача 2. Ограничение параллелизма

Сценарий. Те же запросы, но их 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. Результаты складывайте по индексу, а не по порядку завершения.

Эталон и разбор
final/Practice.cs
        var results = new string[keys.Count];
        await Parallel.ForEachAsync(
            Enumerable.Range(0, keys.Count),
            new ParallelOptions { MaxDegreeOfParallelism = maxParallel, CancellationToken = ct },
            async (i, token) => results[i] = await api.GetAsync(keys[i], token));
        return results;
  • Parallel.ForEachAsync сам держит не больше MaxDegreeOfParallelism задач и перестаёт брать новые после отмены.
  • Запись results[i] = ... в разные ячейки массива из разных потоков безопасна: каждый индекс пишет ровно один обработчик.
  • Вариант на SemaphoreSlim: await sem.WaitAsync(ct) перед запросом и Release() в finally, но тогда все 40 задач уже созданы и ждут у семафора; Parallel.ForEachAsync создаёт их по мере освобождения.

Задача 3. Таймаут и отмена

Сценарий. Внешний сервис иногда «задумывается» на секунды. Нужен запрос с таймаутом; при этом вызывающий может отменить его сам.

Сигнатура: Task<string> GetWithTimeoutAsync(FakeApi api, string key, TimeSpan timeout, CancellationToken ct).

Требования: не уложились в таймаут — TimeoutException, и сам запрос отменён (а не брошен в фоне); отменил вызывающий — OperationCanceledException, а не TimeoutException.

Проверка: запрос на 2 с с таймаутом 200 мс падает с TimeoutException быстро, и у FakeApi Cancelled == 1; быстрый запрос возвращает значение; внешняя отмена через 100 мс даёт OperationCanceledException.

Подсказка

Нужен один токен, который сработает и от внешней отмены, и от таймера (глава 9). Как отличить, кто отменил, когда пришёл OperationCanceledException?

Эталон и разбор
final/Practice.cs
        using var cts = CancellationTokenSource.CreateLinkedTokenSource(ct);
        cts.CancelAfter(timeout);
        try
        {
            return await api.GetAsync(key, cts.Token);
        }
        catch (OperationCanceledException) when (!ct.IsCancellationRequested)
        {
            throw new TimeoutException($"{key}: не уложились в {timeout.TotalMilliseconds:F0} мс");
        }
  • CreateLinkedTokenSource + CancelAfter даёт токен «внешний ИЛИ таймер»; using освобождает связанный токен (глава 9: без Dispose обработчики копятся на родителе).
  • Фильтр when (!ct.IsCancellationRequested) различает причины: если внешний токен не отменён, отменил таймер.
  • Решение WhenAny(запрос, Delay(timeout)) здесь не пройдёт: оно возвращает управление по таймауту, но запрос продолжает работать (глава 9: «таймаут» без отмены). task.WaitAsync(timeout) тоже перестаёт ждать, но не отменяет запрос.

Задача 4. Повтор с экспоненциальной паузой

Сценарий. Сервис иногда отвечает 503. Повторяйте запрос, но не бесконечно и не вплотную.

Сигнатура: Task<string> GetWithRetryAsync(FakeApi api, string key, int maxAttempts, TimeSpan baseDelay, CancellationToken ct).

Требования: повтор только при HttpRequestException; всего не больше maxAttempts обращений; паузы между попытками baseDelay, 2×baseDelay, 4×baseDelay…; после последней неудачи пробрасывается её исключение; отмена прерывает и запрос, и паузу и не считается поводом для повтора.

Проверка: сервис падает дважды и отвечает на третий раз (3 обращения, суммарная пауза ≥ 300 мс); всегда падающий сервис — ровно 3 попытки и HttpRequestException; отмена через 150 мс во время паузы — OperationCanceledException без новых обращений.

Подсказка

Цикл с try/catch вокруг await. Условие when (attempt < maxAttempts) в catch избавляет от отдельной проверки «последняя ли попытка»: на последней исключение просто не перехватывается. Пауза — Task.Delay с токеном.

Эталон и разбор
final/Practice.cs
        for (int attempt = 1; ; attempt++)
        {
            try
            {
                return await api.GetAsync(key, ct);
            }
            catch (HttpRequestException) when (attempt < maxAttempts)
            {
                await Task.Delay(baseDelay * (1 << (attempt - 1)), ct);
            }
        }
  • Ловим только HttpRequestException: OperationCanceledException не перехватывается и сразу выходит наружу.
  • Пауза Task.Delay(..., ct) отменяема. Без токена отмена во время паузы заметилась бы только после неё.
  • На последней попытке фильтр when ложный, исключение летит вызывающему как есть, со своим стеком.
  • В реальном коде добавляют случайный разброс (jitter) к паузе и ограничивают максимальную паузу; готовое решение — библиотека Polly или Microsoft.Extensions.Http.Resilience.

Задача 5. Первый успешный из зеркал

Сценарий. Один и тот же файл есть на трёх зеркалах. Спросите все сразу и возьмите первый успешный ответ; остальные запросы отмените, чтобы не тратить трафик.

Сигнатура: Task<string> FirstSuccessAsync(IReadOnlyList<FakeApi> mirrors, string key, CancellationToken ct).

Требования: быстро упавшее зеркало не должно «выиграть»; результат — первое успешное; остальные запросы отменены; если упали все — AggregateException со всеми ошибками.

Проверка: зеркало 0 падает за 50 мс, зеркало 1 отвечает за 150 мс, зеркало 2 — за 2 с. Ждём ответ зеркала 1 за 150–500 мс и Cancelled == 1 у зеркала 2. Если падают все три, AggregateException с тремя ошибками.

Подсказка

Task.WhenAny вернёт первую завершённую задачу, а она может быть упавшей. Нужен цикл: ждём очередную, смотрим на результат, упавшую убираем из списка, пока не кончатся. Отмена остальных — через общий связанный токен.

Эталон и разбор
final/Practice.cs
        using var cts = CancellationTokenSource.CreateLinkedTokenSource(ct);
        var pending = mirrors.Select(m => m.GetAsync(key, cts.Token)).ToList();
        var errors = new List<Exception>();
        while (pending.Count > 0)
        {
            Task<string> done = await Task.WhenAny(pending);
            pending.Remove(done);
            if (done.IsCompletedSuccessfully)
            {
                cts.Cancel();
                return done.Result;
            }
            errors.Add(done.Exception!.InnerException ?? done.Exception);
        }
        ct.ThrowIfCancellationRequested();
        throw new AggregateException(errors);
  • Цикл WhenAny + Remove для небольшого числа задач — нормально. Для тысяч задач он квадратичен (глава 10: 16 000 задач заняли около 2 с); тогда Task.WhenEach (.NET 9+).
  • cts.Cancel() отменяет все оставшиеся запросы; они завершатся отменой, но это нас больше не касается: возвращённый результат уже взят.
  • Исключения отменённых задач никто не наблюдает, и это безопасно: отменённая задача не попадает в UnobservedTaskException.
  • Строка ct.ThrowIfCancellationRequested() перед throw: если упали все из-за внешней отмены, правильный ответ — отмена, а не AggregateException.

Задача 6. Асинхронный кеш без дублей

Сценарий. Сотня одновременных запросов просит один и тот же справочник. Загружать его надо один раз, остальные ждут тот же результат. Если загрузка упала, ошибку нельзя запоминать навсегда.

Класс: 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.

Эталон и разбор
final/Practice.cs
        private readonly ConcurrentDictionary<string, Lazy<Task<string>>> _items = new();

        public Task<string> GetAsync(string key)
        {
            var lazy = _items.GetOrAdd(key, k =>
            {
                Lazy<Task<string>>? self = null;
                self = new Lazy<Task<string>>(() => LoadAsync(k, self!));
                return self;
            });
            return lazy.Value;
        }

        private async Task<string> LoadAsync(string key, Lazy<Task<string>> self)
        {
            try
            {
                return await factory(key);
            }
            catch
            {
                _items.TryRemove(new KeyValuePair<string, Lazy<Task<string>>>(key, self));   // только свою запись
                throw;
            }
        }
  • В словаре лежит Lazy<Task<string>>: GetOrAdd может создать несколько Lazy, но .Value у проигравших не вызывается, поэтому фабрика запускается один раз.
  • Все ожидающие получают один и тот же Task, то есть одно значение или одну ошибку.
  • При ошибке мы удаляем именно свою запись (TryRemove(KeyValuePair)), чтобы не стереть более новую запись, появившуюся после сбоя.
  • Lazy<Task<T>> без удаления при ошибке кеширует сбой навсегда (глава 10, «Lazy<Task> кеширует ошибку»).
  • Чего здесь нет: отмены. Если один из ожидающих отменится, общую загрузку отменять нельзя, потому что она нужна остальным; поэтому фабрика токена не получает.

Задача 7. Постраничная лента как поток

Сценарий. Лента отдаётся страницами по 10 элементов. Нужен IAsyncEnumerable<string>, из которого потребитель берёт элементы по одному и может остановиться в любой момент.

Сигнатура: IAsyncEnumerable<string> ReadFeedAsync(FakeFeed feed, CancellationToken ct = default).

Требования: страницы читаются по мере надобности (следующую не запрашивать, пока потребитель не дошёл до конца текущей); отмена через WithCancellation(token) прерывает чтение.

Проверка: полный проход — 50 элементов по порядку и ровно 5 запросов страниц; потребитель остановился на 15-м элементе — запрошено не больше 2 страниц; WithCancellation с токеном, отменяемым через 80 мс, даёт OperationCanceledException.

Подсказка

Асинхронный итератор: async IAsyncEnumerable<T> с yield return (глава 11). Чтобы токен из WithCancellation попал в ваш метод, параметр нужно пометить атрибутом [EnumeratorCancellation].

Эталон и разбор
final/Practice.cs
        for (int page = 0; ; page++)
        {
            Page p = await feed.GetPageAsync(page, ct);
            foreach (var item in p.Items) yield return item;
            if (!p.HasNext) yield break;
        }
  • Итератор ленив: код после yield return выполняется только когда потребитель запрашивает следующий элемент, поэтому «читать вперёд» не получается само собой.
  • Без [EnumeratorCancellation] токен, переданный через WithCancellation, до вашего ct не дойдёт, и отмена будет работать только между элементами.
  • HasNext избавляет от лишнего запроса пустой страницы в конце.
  • Если нужно предзагружать следующую страницу в фоне (быстрее, но тратит запросы), это уже Channel с ёмкостью 1: см. задачу 8.

Задача 8. Конвейер на Channel

Сценарий. Нужно обработать 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, когда все читатели закончили.

Эталон и разбор
final/Practice.cs
        using var cts = CancellationTokenSource.CreateLinkedTokenSource(ct);
        var input = Channel.CreateBounded<string>(20);
        var fetched = Channel.CreateBounded<string>(20);

        async Task Guard(Func<Task> stage)
        {
            try { await stage(); }
            catch { cts.Cancel(); throw; }      // ошибка одной стадии останавливает остальные
        }

        var producer = Guard(async () =>
        {
            foreach (var key in keys) await input.Writer.WriteAsync(key, cts.Token);
            input.Writer.Complete();
        });

        var fetchers = Enumerable.Range(0, workers).Select(_ => Guard(async () =>
        {
            await foreach (var key in input.Reader.ReadAllAsync(cts.Token))
                await fetched.Writer.WriteAsync(await api.GetAsync(key, cts.Token), cts.Token);
        })).ToArray();

        var closer = Guard(async () =>
        {
            await Task.WhenAll(fetchers);
            fetched.Writer.Complete();
        });

        int total = 0;
        var writer = Guard(async () =>
        {
            var batch = new List<string>(batchSize);
            await foreach (var row in fetched.Reader.ReadAllAsync(cts.Token))
            {
                batch.Add(row);
                if (batch.Count < batchSize) continue;
                await db.InsertBatchAsync(batch.ToArray(), cts.Token);
                total += batch.Count;
                batch.Clear();
            }
            if (batch.Count > 0)
            {
                await db.InsertBatchAsync(batch.ToArray(), cts.Token);
                total += batch.Count;
            }
        });

        var all = Task.WhenAll(producer, closer, writer);
        try
        {
            await all;
        }
        catch
        {
            // WhenAll бросит первое по порядку исключение, а оно может быть OperationCanceledException от «остановленной» стадии.
            var real = all.Exception!.InnerExceptions.FirstOrDefault(e => e is not OperationCanceledException) ?? all.Exception.InnerExceptions[0];
            ExceptionDispatchInfo.Capture(real).Throw();
        }
        return total;
  • 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 для запросов и пакеты в потокобезопасный буфер. Он закрывает ограничение запросов, но обратное давление от медленной базы и единственного писателя придётся добавлять отдельно.

20.3. Типичные ошибки, которые ловят проверки

Ошибка Какая проверка падает
await в цикле вместо WhenAll 1: «10 запросов заняли ~1000 мс»
WhenAll по всем ключам без ограничения 2: «одновременно шло 40 запросов»
Таймаут через WhenAny + Delay, запрос не отменён 3: «сам запрос не отменён»
Пауза между попытками без роста или Delay без токена 4: «паузы слишком короткие» / отмена не прерывает паузу
WhenAny без проверки, что задача не упала 5: падает исключение быстрого зеркала
GetOrAdd(key, factory) без Lazy или ошибка осталась в кеше 6: фабрика вызвана несколько раз / «упавшая загрузка осталась в кеше»
Все страницы загружаются заранее 7: «лента читает вперёд»
Нет общего токена и Guard 8: конвейер не останавливается после ошибки

Итоги

  • Типовая прикладная задача состоит из пяти кирпичей: WhenAll, ограничение параллелизма, связанный токен с таймаутом, цикл повтора, Channel.
  • В каждом кирпиче есть ловушка, которую не показывает «зелёный» тест на счастливом пути: потеря отмены, кеширование ошибки, не та причина в исключении.
  • Эталон проходит проверки, но реальному коду нужны ещё логирование, метрики и политики повторов из готовых библиотек.

Код лабораторной

Запуск из папки главы: dotnet run -c Release --project start (все задачи) или -- 3 5 (выбранные).

start/Practice.cs
using System.Runtime.CompilerServices;

// Восемь задач. Тела методов пустые (NotImplementedException): напишите их сами. Эталонные решения — в final/. Описание задач и разбор — в README главы.
static class Practice
{
    /// <summary>
    /// Задача 1. Получить данные по всем ключам <b>параллельно</b>. Порядок результатов — как порядок ключей.
    /// Если любой запрос упал — исключение этого запроса.
    /// </summary>
    public static Task<string[]> FetchAllAsync(FakeApi api, IReadOnlyList<string> keys, CancellationToken ct)
    {
        // TODO: ваше решение
        throw new NotImplementedException();
    }

    /// <summary>
    /// Задача 2. То же, но одновременно выполняется не больше <paramref name="maxParallel"/> запросов.
    /// Порядок результатов — как порядок ключей. При отмене новые запросы не запускаются.
    /// </summary>
    public static async Task<string[]> FetchLimitedAsync(FakeApi api, IReadOnlyList<string> keys, int maxParallel, CancellationToken ct)
    {
        // TODO: ваше решение
        throw new NotImplementedException();
    }

    /// <summary>
    /// Задача 3. Запрос с таймаутом. Не уложились — <see cref="TimeoutException"/>, и сам запрос должен быть отменён.
    /// Если отменил вызывающий (<paramref name="ct"/>) — <see cref="OperationCanceledException"/>, а не TimeoutException.
    /// </summary>
    public static async Task<string> GetWithTimeoutAsync(FakeApi api, string key, TimeSpan timeout, CancellationToken ct)
    {
        // TODO: ваше решение
        throw new NotImplementedException();
    }

    /// <summary>
    /// Задача 4. Повтор при <see cref="HttpRequestException"/>: до <paramref name="maxAttempts"/> обращений,
    /// между попытками пауза baseDelay, 2×baseDelay, 4×baseDelay… Последняя ошибка пробрасывается.
    /// Отмена прерывает и запрос, и паузу и повторов не вызывает.
    /// </summary>
    public static async Task<string> GetWithRetryAsync(FakeApi api, string key, int maxAttempts, TimeSpan baseDelay, CancellationToken ct)
    {
        // TODO: ваше решение
        throw new NotImplementedException();
    }

    /// <summary>
    /// Задача 5. Одновременно спросить все зеркала и вернуть ответ первого <b>успешного</b>. Остальные запросы отменить.
    /// Если упали все — <see cref="AggregateException"/> со всеми ошибками.
    /// </summary>
    public static async Task<string> FirstSuccessAsync(IReadOnlyList<FakeApi> mirrors, string key, CancellationToken ct)
    {
        // TODO: ваше решение
        throw new NotImplementedException();
    }

    /// <summary>
    /// Задача 6. Асинхронный кеш. Одновременные запросы одного ключа вызывают фабрику <b>один раз</b>;
    /// готовое значение хранится. Упавшая загрузка в кеше не остаётся: следующий запрос пробует заново.
    /// </summary>
    public sealed class AsyncCache(Func<string, Task<string>> factory)
    {
        public Task<string> GetAsync(string key)
        {
            // TODO: ваше решение (factory(key) вызывается в вашем коде)
            throw new NotImplementedException();
        }
    }

    /// <summary>
    /// Задача 7. Лента как <see cref="IAsyncEnumerable{T}"/>: страницы читаются по мере надобности (страницу N+1 не запрашивать,
    /// пока потребитель не дошёл до конца страницы N). Отмена через WithCancellation должна прерывать чтение.
    /// </summary>
    public static async IAsyncEnumerable<string> ReadFeedAsync(FakeFeed feed, [EnumeratorCancellation] CancellationToken ct = default)
    {
        // TODO: ваше решение (в теле метода должен быть yield return)
        throw new NotImplementedException();
#pragma warning disable CS0162
        yield break;   // заглушка, чтобы метод компилировался; удалите вместе с pragma
#pragma warning restore CS0162
    }

    /// <summary>
    /// Задача 8. Конвейер «ключи → API → база»: <paramref name="workers"/> параллельных запросов к API,
    /// запись в базу одним писателем пакетами по <paramref name="batchSize"/> (последний пакет может быть меньше).
    /// Очереди между стадиями ограничены (обратное давление). Ошибка любой стадии отменяет остальные и
    /// пробрасывается наружу (настоящая причина, а не OperationCanceledException). Возвращает число записанных строк.
    /// </summary>
    public static async Task<int> RunPipelineAsync(FakeApi api, FakeDb db, IEnumerable<string> keys, int workers, int batchSize, CancellationToken ct)
    {
        // TODO: ваше решение
        throw new NotImplementedException();
    }
}
final/Practice.cs
using System.Collections.Concurrent;
using System.Runtime.CompilerServices;
using System.Runtime.ExceptionServices;
using System.Threading.Channels;

// Восемь задач. В start/ тела методов пустые (NotImplementedException): напишите их сами.
// В final/ лежат эталонные решения. Описание задач и разбор — в README главы.
static class Practice
{
    /// <summary>
    /// Задача 1. Получить данные по всем ключам <b>параллельно</b>. Порядок результатов — как порядок ключей.
    /// Если любой запрос упал — исключение этого запроса.
    /// </summary>
    public static Task<string[]> FetchAllAsync(FakeApi api, IReadOnlyList<string> keys, CancellationToken ct)
    {
        // <решение>
        return Task.WhenAll(keys.Select(k => api.GetAsync(k, ct)));
        // </решение>
    }

    /// <summary>
    /// Задача 2. То же, но одновременно выполняется не больше <paramref name="maxParallel"/> запросов.
    /// Порядок результатов — как порядок ключей. При отмене новые запросы не запускаются.
    /// </summary>
    public static async Task<string[]> FetchLimitedAsync(FakeApi api, IReadOnlyList<string> keys, int maxParallel, CancellationToken ct)
    {
        // <решение>
        var results = new string[keys.Count];
        await Parallel.ForEachAsync(
            Enumerable.Range(0, keys.Count),
            new ParallelOptions { MaxDegreeOfParallelism = maxParallel, CancellationToken = ct },
            async (i, token) => results[i] = await api.GetAsync(keys[i], token));
        return results;
        // </решение>
    }

    /// <summary>
    /// Задача 3. Запрос с таймаутом. Не уложились — <see cref="TimeoutException"/>, и сам запрос должен быть отменён.
    /// Если отменил вызывающий (<paramref name="ct"/>) — <see cref="OperationCanceledException"/>, а не TimeoutException.
    /// </summary>
    public static async Task<string> GetWithTimeoutAsync(FakeApi api, string key, TimeSpan timeout, CancellationToken ct)
    {
        // <решение>
        using var cts = CancellationTokenSource.CreateLinkedTokenSource(ct);
        cts.CancelAfter(timeout);
        try
        {
            return await api.GetAsync(key, cts.Token);
        }
        catch (OperationCanceledException) when (!ct.IsCancellationRequested)
        {
            throw new TimeoutException($"{key}: не уложились в {timeout.TotalMilliseconds:F0} мс");
        }
        // </решение>
    }

    /// <summary>
    /// Задача 4. Повтор при <see cref="HttpRequestException"/>: до <paramref name="maxAttempts"/> обращений,
    /// между попытками пауза baseDelay, 2×baseDelay, 4×baseDelay… Последняя ошибка пробрасывается.
    /// Отмена прерывает и запрос, и паузу и повторов не вызывает.
    /// </summary>
    public static async Task<string> GetWithRetryAsync(FakeApi api, string key, int maxAttempts, TimeSpan baseDelay, CancellationToken ct)
    {
        // <решение>
        for (int attempt = 1; ; attempt++)
        {
            try
            {
                return await api.GetAsync(key, ct);
            }
            catch (HttpRequestException) when (attempt < maxAttempts)
            {
                await Task.Delay(baseDelay * (1 << (attempt - 1)), ct);
            }
        }
        // </решение>
    }

    /// <summary>
    /// Задача 5. Одновременно спросить все зеркала и вернуть ответ первого <b>успешного</b>. Остальные запросы отменить.
    /// Если упали все — <see cref="AggregateException"/> со всеми ошибками.
    /// </summary>
    public static async Task<string> FirstSuccessAsync(IReadOnlyList<FakeApi> mirrors, string key, CancellationToken ct)
    {
        // <решение>
        using var cts = CancellationTokenSource.CreateLinkedTokenSource(ct);
        var pending = mirrors.Select(m => m.GetAsync(key, cts.Token)).ToList();
        var errors = new List<Exception>();
        while (pending.Count > 0)
        {
            Task<string> done = await Task.WhenAny(pending);
            pending.Remove(done);
            if (done.IsCompletedSuccessfully)
            {
                cts.Cancel();
                return done.Result;
            }
            errors.Add(done.Exception!.InnerException ?? done.Exception);
        }
        ct.ThrowIfCancellationRequested();
        throw new AggregateException(errors);
        // </решение>
    }

    /// <summary>
    /// Задача 6. Асинхронный кеш. Одновременные запросы одного ключа вызывают фабрику <b>один раз</b>;
    /// готовое значение хранится. Упавшая загрузка в кеше не остаётся: следующий запрос пробует заново.
    /// </summary>
    public sealed class AsyncCache(Func<string, Task<string>> factory)
    {
        // <решение>
        private readonly ConcurrentDictionary<string, Lazy<Task<string>>> _items = new();

        public Task<string> GetAsync(string key)
        {
            var lazy = _items.GetOrAdd(key, k =>
            {
                Lazy<Task<string>>? self = null;
                self = new Lazy<Task<string>>(() => LoadAsync(k, self!));
                return self;
            });
            return lazy.Value;
        }

        private async Task<string> LoadAsync(string key, Lazy<Task<string>> self)
        {
            try
            {
                return await factory(key);
            }
            catch
            {
                _items.TryRemove(new KeyValuePair<string, Lazy<Task<string>>>(key, self));   // только свою запись
                throw;
            }
        }
        // </решение>
    }

    /// <summary>
    /// Задача 7. Лента как <see cref="IAsyncEnumerable{T}"/>: страницы читаются по мере надобности (страницу N+1 не запрашивать,
    /// пока потребитель не дошёл до конца страницы N). Отмена через WithCancellation должна прерывать чтение.
    /// </summary>
    public static async IAsyncEnumerable<string> ReadFeedAsync(FakeFeed feed, [EnumeratorCancellation] CancellationToken ct = default)
    {
        // <решение>
        for (int page = 0; ; page++)
        {
            Page p = await feed.GetPageAsync(page, ct);
            foreach (var item in p.Items) yield return item;
            if (!p.HasNext) yield break;
        }
        // </решение>
    }

    /// <summary>
    /// Задача 8. Конвейер «ключи → API → база»: <paramref name="workers"/> параллельных запросов к API,
    /// запись в базу одним писателем пакетами по <paramref name="batchSize"/> (последний пакет может быть меньше).
    /// Очереди между стадиями ограничены (обратное давление). Ошибка любой стадии отменяет остальные и
    /// пробрасывается наружу (настоящая причина, а не OperationCanceledException). Возвращает число записанных строк.
    /// </summary>
    public static async Task<int> RunPipelineAsync(FakeApi api, FakeDb db, IEnumerable<string> keys, int workers, int batchSize, CancellationToken ct)
    {
        // <решение>
        using var cts = CancellationTokenSource.CreateLinkedTokenSource(ct);
        var input = Channel.CreateBounded<string>(20);
        var fetched = Channel.CreateBounded<string>(20);

        async Task Guard(Func<Task> stage)
        {
            try { await stage(); }
            catch { cts.Cancel(); throw; }      // ошибка одной стадии останавливает остальные
        }

        var producer = Guard(async () =>
        {
            foreach (var key in keys) await input.Writer.WriteAsync(key, cts.Token);
            input.Writer.Complete();
        });

        var fetchers = Enumerable.Range(0, workers).Select(_ => Guard(async () =>
        {
            await foreach (var key in input.Reader.ReadAllAsync(cts.Token))
                await fetched.Writer.WriteAsync(await api.GetAsync(key, cts.Token), cts.Token);
        })).ToArray();

        var closer = Guard(async () =>
        {
            await Task.WhenAll(fetchers);
            fetched.Writer.Complete();
        });

        int total = 0;
        var writer = Guard(async () =>
        {
            var batch = new List<string>(batchSize);
            await foreach (var row in fetched.Reader.ReadAllAsync(cts.Token))
            {
                batch.Add(row);
                if (batch.Count < batchSize) continue;
                await db.InsertBatchAsync(batch.ToArray(), cts.Token);
                total += batch.Count;
                batch.Clear();
            }
            if (batch.Count > 0)
            {
                await db.InsertBatchAsync(batch.ToArray(), cts.Token);
                total += batch.Count;
            }
        });

        var all = Task.WhenAll(producer, closer, writer);
        try
        {
            await all;
        }
        catch
        {
            // WhenAll бросит первое по порядку исключение, а оно может быть OperationCanceledException от «остановленной» стадии.
            var real = all.Exception!.InnerExceptions.FirstOrDefault(e => e is not OperationCanceledException) ?? all.Exception.InnerExceptions[0];
            ExceptionDispatchInfo.Capture(real).Throw();
        }
        return total;
        // </решение>
    }
}
start/Fakes.cs
using System.Collections.Concurrent;

// «Внешний мир» для задач: ни сети, ни базы не нужно, но задержки, сбои и отмена ведут себя как настоящие.
// Эти файлы трогать не надо: решайте в Practice.cs.

/// <summary>Условный HTTP-сервис. Считает вызовы, одновременные вызовы и отменённые вызовы.</summary>
sealed class FakeApi(int latencyMs = 100)
{
    private int _inFlight, _peak, _calls, _cancelled;
    private readonly ConcurrentDictionary<string, int> _attempts = new();

    public int LatencyMs { get; set; } = latencyMs;
    /// <summary>Первые N обращений к каждому ключу падают с HttpRequestException (503).</summary>
    public int FailFirstAttempts { get; set; }
    public bool AlwaysFail { get; set; }
    public HashSet<string> BadKeys { get; } = [];

    public int Calls => Volatile.Read(ref _calls);
    public int Peak => Volatile.Read(ref _peak);
    public int Cancelled => Volatile.Read(ref _cancelled);

    public async Task<string> GetAsync(string key, CancellationToken ct = default)
    {
        Interlocked.Increment(ref _calls);
        int now = Interlocked.Increment(ref _inFlight);
        int seen;
        while (now > (seen = Volatile.Read(ref _peak)) && Interlocked.CompareExchange(ref _peak, now, seen) != seen) { }
        try
        {
            await Task.Delay(LatencyMs, ct);
            int attempt = _attempts.AddOrUpdate(key, 1, (_, n) => n + 1);
            if (AlwaysFail || BadKeys.Contains(key) || attempt <= FailFirstAttempts)
                throw new HttpRequestException($"503 для {key} (попытка {attempt})");
            return $"data:{key}";
        }
        catch (OperationCanceledException)
        {
            Interlocked.Increment(ref _cancelled);
            throw;
        }
        finally
        {
            Interlocked.Decrement(ref _inFlight);
        }
    }
}

/// <summary>Условная база данных: пакетная вставка. Запоминает всё записанное.</summary>
sealed class FakeDb
{
    private int _writing, _peakWriting;
    public ConcurrentBag<string> Rows { get; } = [];
    public ConcurrentQueue<int> BatchSizes { get; } = [];
    public int PeakConcurrentWrites => Volatile.Read(ref _peakWriting);

    public async Task InsertBatchAsync(IReadOnlyCollection<string> rows, CancellationToken ct = default)
    {
        int now = Interlocked.Increment(ref _writing);
        int seen;
        while (now > (seen = Volatile.Read(ref _peakWriting)) && Interlocked.CompareExchange(ref _peakWriting, now, seen) != seen) { }
        try
        {
            await Task.Delay(30, ct);
            BatchSizes.Enqueue(rows.Count);
            foreach (var r in rows) Rows.Add(r);
        }
        finally { Interlocked.Decrement(ref _writing); }
    }
}

sealed record Page(IReadOnlyList<string> Items, bool HasNext);

/// <summary>Постраничная лента: 5 страниц по 10 элементов.</summary>
sealed class FakeFeed
{
    private int _pages;
    public int PagesRequested => Volatile.Read(ref _pages);

    public async Task<Page> GetPageAsync(int page, CancellationToken ct = default)
    {
        Interlocked.Increment(ref _pages);
        await Task.Delay(50, ct);
        var items = Enumerable.Range(page * 10, 10).Select(i => $"item{i}").ToList();
        return new Page(items, page < 4);
    }
}
start/Checks.cs
using System.Diagnostics;

// Проверки для самопроверки. Каждая запускает ваше решение из Practice.cs на «внешнем мире» из Fakes.cs.
static class Checks
{
    public static async Task RunAsync(string[] args)
    {
        var all = new (int n, string title, Func<Task<string?>> run)[]
        {
            (1, "Параллельный сбор", Check1), (2, "Ограничение параллелизма", Check2),
            (3, "Таймаут и отмена", Check3), (4, "Повтор с экспоненциальной паузой", Check4),
            (5, "Первый успешный из зеркал", Check5), (6, "Кеш без дублей", Check6),
            (7, "Постраничная лента как поток", Check7), (8, "Конвейер на Channel", Check8),
        };
        int passed = 0, ran = 0;
        foreach (var (n, title, run) in all.Where(t => args.Length == 0 || args.Contains(t.n.ToString())))
        {
            ran++;
            var sw = Stopwatch.StartNew();
            string? error;
            try
            {
                var task = run();
                error = await Task.WhenAny(task, Task.Delay(15_000)) == task ? await task : "не уложились в 15 секунд (зависло?)";
            }
            catch (NotImplementedException) { error = "не реализовано"; }
            catch (Exception e) { error = $"исключение {e.GetType().Name}: {e.Message}"; }
            if (error is null) passed++;
            Console.WriteLine($"{(error is null ? "PASS" : "FAIL")}  {n}. {title,-34} {sw.ElapsedMilliseconds,5} мс{(error is null ? "" : "   → " + error)}");
        }
        Console.WriteLine($"\nИтого: {passed} из {ran}");
    }

    static string[] Keys(int n) => Enumerable.Range(0, n).Select(i => $"k{i}").ToArray();
    static string[] Expected(IEnumerable<string> keys) => keys.Select(k => $"data:{k}").ToArray();

    // ---- 1 ----
    static async Task<string?> Check1()
    {
        var api = new FakeApi(100);
        var keys = Keys(10);
        var sw = Stopwatch.StartNew();
        var r = await Practice.FetchAllAsync(api, keys, CancellationToken.None);
        if (!r.SequenceEqual(Expected(keys))) return "результаты не совпали или порядок нарушен";
        if (sw.ElapsedMilliseconds > 450) return $"10 запросов по 100 мс заняли {sw.ElapsedMilliseconds} мс: они идут не параллельно";
        return null;
    }

    // ---- 2 ----
    static async Task<string?> Check2()
    {
        var api = new FakeApi(50);
        var keys = Keys(40);
        var sw = Stopwatch.StartNew();
        var r = await Practice.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)";

        var api2 = new FakeApi(100);
        using var cts = new CancellationTokenSource(250);
        try { await Practice.FetchLimitedAsync(api2, Keys(100), 5, cts.Token); return "при отмене не было исключения"; }
        catch (OperationCanceledException) { }
        if (api2.Calls > 20) return $"после отмены продолжали запускаться запросы: всего {api2.Calls}";
        return null;
    }

    // ---- 3 ----
    static async Task<string?> Check3()
    {
        var slow = new FakeApi(2000);
        var sw = Stopwatch.StartNew();
        try { await Practice.GetWithTimeoutAsync(slow, "a", TimeSpan.FromMilliseconds(200), CancellationToken.None); return "таймаут не сработал"; }
        catch (TimeoutException) { }
        if (sw.ElapsedMilliseconds > 700) return $"таймаут 200 мс сработал через {sw.ElapsedMilliseconds} мс";
        await Task.Delay(100);
        if (slow.Cancelled != 1) return $"сам запрос не отменён (отменено {slow.Cancelled}): после таймаута он продолжает работать впустую";

        var fast = new FakeApi(50);
        if (await Practice.GetWithTimeoutAsync(fast, "b", TimeSpan.FromSeconds(1), CancellationToken.None) != "data:b") return "быстрый запрос вернул не то";

        using var cts = new CancellationTokenSource(100);
        try { await Practice.GetWithTimeoutAsync(new FakeApi(2000), "c", TimeSpan.FromSeconds(5), cts.Token); return "внешняя отмена не сработала"; }
        catch (TimeoutException) { return "внешняя отмена превратилась в TimeoutException: различайте причины"; }
        catch (OperationCanceledException) { }
        return null;
    }

    // ---- 4 ----
    static async Task<string?> Check4()
    {
        var api = new FakeApi(10) { FailFirstAttempts = 2 };
        var sw = Stopwatch.StartNew();
        var r = await Practice.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} мс";

        var dead = new FakeApi(10) { AlwaysFail = true };
        try { await Practice.GetWithRetryAsync(dead, "y", 3, TimeSpan.FromMilliseconds(20), CancellationToken.None); return "ошибка проглочена"; }
        catch (HttpRequestException) { }
        if (dead.Calls != 3) return $"ожидали ровно 3 попытки, было {dead.Calls}";

        var api3 = new FakeApi(10) { AlwaysFail = true };
        using var cts = new CancellationTokenSource(150);
        try { await Practice.GetWithRetryAsync(api3, "z", 10, TimeSpan.FromMilliseconds(100), cts.Token); return "при отмене не было исключения"; }
        catch (OperationCanceledException) { }
        catch (HttpRequestException) { return "отмена во время паузы превратилась в HttpRequestException или повторы не прерываются"; }
        if (api3.Calls > 3) return $"после отмены повторы продолжались: {api3.Calls} обращений";
        return null;
    }

    // ---- 5 ----
    static async Task<string?> Check5()
    {
        var m0 = new FakeApi(50) { AlwaysFail = true };
        var m1 = new FakeApi(150);
        var m2 = new FakeApi(2000);
        var sw = Stopwatch.StartNew();
        var r = await Practice.FirstSuccessAsync([m0, m1, m2], "k", CancellationToken.None);
        if (r != "data:k") return "результат не тот";
        if (sw.ElapsedMilliseconds > 500) return $"ждали {sw.ElapsedMilliseconds} мс: должен вернуться первый успешный (около 150 мс)";
        await Task.Delay(100);
        if (m2.Cancelled != 1) return "медленное зеркало не отменено: запрос продолжает работать впустую";

        try
        {
            await Practice.FirstSuccessAsync([new FakeApi(20) { AlwaysFail = true }, new FakeApi(30) { AlwaysFail = true }, new FakeApi(40) { AlwaysFail = true }], "k", CancellationToken.None);
            return "когда падают все, нужно исключение";
        }
        catch (AggregateException e) when (e.InnerExceptions.Count == 3) { }
        catch (AggregateException e) { return $"AggregateException должно содержать все 3 ошибки, а в нём {e.InnerExceptions.Count}"; }
        return null;
    }

    // ---- 6 ----
    static async Task<string?> Check6()
    {
        int calls = 0;
        var cache = new Practice.AsyncCache(async key => { Interlocked.Increment(ref calls); await Task.Delay(100); return $"v:{key}"; });
        var results = await Task.WhenAll(Enumerable.Range(0, 20).Select(_ => cache.GetAsync("a")));
        if (results.Any(x => x != "v:a")) return "значения не совпали";
        if (calls != 1) return $"20 одновременных запросов одного ключа вызвали фабрику {calls} раз (нужно 1)";
        await cache.GetAsync("a");
        if (calls != 1) return "повторный запрос готового ключа снова вызвал фабрику";
        await cache.GetAsync("b");
        if (calls != 2) return "другой ключ должен вызывать фабрику отдельно";

        int attempts = 0;
        var flaky = new Practice.AsyncCache(async key =>
        {
            int n = Interlocked.Increment(ref attempts);
            await Task.Delay(50);
            if (n == 1) throw new InvalidOperationException("первый вызов падает");
            return "ok";
        });
        try { await flaky.GetAsync("x"); return "первый вызов должен был упасть"; }
        catch (InvalidOperationException) { }
        try
        {
            if (await flaky.GetAsync("x") != "ok") return "после сбоя вернулось не то значение";
        }
        catch (InvalidOperationException) { return "упавшая загрузка осталась в кеше: следующий вызов получил ту же ошибку (нужно пробовать заново)"; }
        return null;
    }

    // ---- 7 ----
    static async Task<string?> Check7()
    {
        var feed = new FakeFeed();
        var all = new List<string>();
        await foreach (var x in Practice.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}";

        var feed2 = new FakeFeed();
        var firstFifteen = new List<string>();
        await foreach (var x in Practice.ReadFeedAsync(feed2))
        {
            firstFifteen.Add(x);
            if (firstFifteen.Count == 15) break;
        }
        await Task.Delay(150);
        if (feed2.PagesRequested > 2) return $"потребитель остановился на 15-м элементе, а страниц запрошено {feed2.PagesRequested} (нужно 2): лента читает вперёд";

        var feed3 = new FakeFeed();
        using var cts = new CancellationTokenSource(80);
        try
        {
            await foreach (var _ in Practice.ReadFeedAsync(feed3).WithCancellation(cts.Token)) { }
            return "WithCancellation(token) не прервал чтение: нужен [EnumeratorCancellation]";
        }
        catch (OperationCanceledException) { }
        return null;
    }

    // ---- 8 ----
    static async Task<string?> Check8()
    {
        var api = new FakeApi(20);
        var db = new FakeDb();
        var keys = Keys(95);
        var sw = Stopwatch.StartNew();
        int n = await Practice.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} потока(ов) одновременно; ожидали одного писателя";
        var sizes = 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} мс";

        var api2 = new FakeApi(20);
        api2.BadKeys.Add("k30");
        var db2 = new FakeDb();
        var sw2 = Stopwatch.StartNew();
        try { await Practice.RunPipelineAsync(api2, db2, Keys(95), 4, 10, CancellationToken.None); return "ошибка на одном ключе должна ломать конвейер"; }
        catch (HttpRequestException) { }
        catch (Exception e) { return $"ожидали HttpRequestException (настоящую причину), получили {e.GetType().Name}"; }
        if (sw2.ElapsedMilliseconds > 3000) return "после ошибки конвейер завершался слишком долго";
        if (api2.Calls >= 95) return "после ошибки конвейер продолжал обрабатывать все ключи: остановите остальных";
        return null;
    }
}
start/Program.cs
1
2
3
4
5
// Глава 20. Практикум: восемь задач на async/await.
// Решайте в Practice.cs, проверяйте:
//   dotnet run -c Release --project start            — все задачи
//   dotnet run -c Release --project start -- 3 5     — только задачи 3 и 5
await Checks.RunAsync(args);