Цель: понять, чем yield return отличается от await, как устроен асинхронный итератор (IAsyncEnumerable<T>, await foreach) и почему он не аллоцирует на элемент. Разобрать по коду Channel<T>: режимы переполнения, обратное давление, завершение, синхронные продолжения. Узнать, что LINQ для IAsyncEnumerable теперь входит в BCL (.NET 10).
Лабораторная:start/ — асинхронный итератор и ограниченный канал с тремя TODO (finally, DropOldest, LINQ). final/ — девять опытов: обычный итератор, await foreach (отмена, break, WithCancellation, AsyncLocal), аллокации, LINQ из BCL, режимы FullMode, завершение канала, несколько читателей, AllowSynchronousContinuations, стоимость элемента канала. Код — в конце главы.
Визуализация:Обратное давление — Играйте производителем и потребителем при каждом FullMode.
Статус: ✅ проверено на Ubuntu 26.04 (2 ядра), рантайм 10.0.12. Листинги CoreLib — декомпиляция System.Private.CoreLib 10.0.12 (.\tools\disasm.ps1 -Assembly corelib -Type …); Channel<T> — из System.Threading.Channels.dll того же рантайма; листинги итераторов — ilspycmd с -ds AsyncAwait=false -ds YieldReturn=false на final. Перепроверено на Windows 11 (8 ядер): вывод совпадает, отличаются только времена и шум замеров аллокаций.
publicstaticIEnumerable<int>Numbers(intmax){Console.WriteLine(" [итератор] тело началось");for(inti=1;i<=max;i++)yieldreturni;Console.WriteLine(" [итератор] тело закончилось");}
yield return x значит «отдать x вызывающему и замереть на этой строке до следующего запроса». Компилятор превращает такой метод в класс (sealed class, не структуру), который реализует IEnumerable<int>иIEnumerator<int> сразу:
Numbers после компиляции (ILSpy без восстановления yield, Release)
privatesealedclass<Numbers>d__0:IEnumerable<int>,IEnumerable,IEnumerator<int>,IEnumerator,IDisposable{privateint<>1__state;privateint<>2__current;privateint<>l__initialThreadId;privateintmax;privateint<i>5__2;public<Numbers>d__0(int<>1__state){this.<>1__state=<>1__state;<>l__initialThreadId=Environment.CurrentManagedThreadId;}privateboolMoveNext(){switch(<>1__state){default:returnfalse;case0:<>1__state=-1;Console.WriteLine(" [итератор] тело началось");<i>5__2=1;break;case1:<>1__state=-1;<i>5__2++;break;}if(<i>5__2<=max){<>2__current=<i>5__2;<>1__state=1;returntrue;}Console.WriteLine(" [итератор] тело закончилось");returnfalse;}IEnumerator<int>IEnumerable<int>.GetEnumerator(){<Numbers>d__0obj;if(<>1__state==-2&&<>l__initialThreadId==Environment.CurrentManagedThreadId){<>1__state=0;obj=this;}else{obj=new<Numbers>d__0(0);}...}}
Строки 3–4 и 7: состояние, текущий элемент (<>2__current) и поднятая локальная i. Те же приёмы, что в главе 2: номер состояния, поля вместо локальных.
Строки 15–39: MoveNext — обычный метод, возвращающий bool. Каждый вызов проходит тело до следующего yield return: кладёт значение в <>2__current, запоминает состояние (1) и возвращает true. Когда тело кончилось, возвращает false. Пока никто не вызвал MoveNext, ни одна строка тела не выполняется: итератор ленив.
Строки 41–54: GetEnumerator. Если последовательность ещё не начата (состояние -2) и вызов сделан на том же потоке, в котором её создали, перечислитель — это сам объект последовательности (строка 47: obj = this). Для второго вызова (или с другого потока) создаётся новый экземпляр. Поэтому на один обход приходится один объект.
== 1. Обычный итератор ==
метод вызван, тело ещё не началось
[итератор] тело началось
получено 1
получено 2
[итератор] тело закончилось
первый GetEnumerator вернул саму последовательность: True, второй: False
тип: ‹Numbers›d__0
Строка 2 вывода:Numbers(2) вызван (строка 13 кода), а «тело началось» напечатано только в первом foreach (строки 3–6).
Строка 7: первый GetEnumerator вернул саму последовательность, второй — новый объект.
yield return (итератор)
await (async-метод)
Что значит пауза
отдать значение вызывающему
ждать результат чего-то
Кто возобновляет
вызывающий, следующим MoveNext(), синхронно, в том же потоке
завершение того, чего ждали (в любом потоке)
Возвращает вызывающему
элементы по одному
Task, один раз в конце
Асинхронный итератор объединяет обе возможности: внутри можно и await, и yield return.
awaitforeach(intninCountAsync(10,ct)){// тело цикла}
Компилятор разворачивает его так (смысл; в машине состояний потребителя это try/catch с запоминанием исключения, потому что в сгенерированном коде await в finally недоступен):
IAsyncEnumerator<int>e=CountAsync(10,ct).GetAsyncEnumerator();try{while(awaite.MoveNextAsync())// ValueTask<bool>{intn=e.Current;// тело цикла}}finally{awaite.DisposeAsync();// и при break, и при исключении}
Строка 4:MoveNextAsync() возвращает ValueTask<bool>: true — есть следующий элемент (он в Current), false — конец.
Строки 10–12: DisposeAsync вызывается всегда, в том числе при break и исключении в теле цикла. В итераторе это выполнит блок finally.
== 2. await foreach ==
получено 1
получено 2
получено 3
получено 4
[итератор] finally
отменено по таймауту: TaskCanceledException
--- break на втором элементе:
получено 1
получено 2
[итератор] finally
--- WithCancellation:
получено 1
получено 2
[итератор] finally
отменено токеном из WithCancellation
первый GetAsyncEnumerator вернул саму последовательность: True, второй: False
--- AsyncLocal, записанный внутри итератора:
итератор без await:
[итератор] до записи Ambient = 1
[итератор] перед yield Ambient = 99
потребитель, элемент 0: Ambient = 1
[итератор] перед yield Ambient = 1
потребитель, элемент 1: Ambient = 1
итератор с await Task.Yield():
[итератор] до записи Ambient = 1
[итератор] перед yield Ambient = 99
потребитель, элемент 0: Ambient = 1
[итератор] перед yield Ambient = 1
потребитель, элемент 1: Ambient = 1
Строки 2–7 вывода (строки 25–33 кода): таймаут 450 мс прервал цикл на пятом шаге: 4 элемента, затем Task.Delay внутри итератора бросил TaskCanceledException (глава 9, §9.2). Блок finally итератора выполнился.
Строки 8–11 (строки 34–39): break на втором элементе. Потребитель вышел из цикла, DisposeAsync разбудил итератор в режиме «закрытия», и finally выполнился. Если бы DisposeAsync не вызвали, finally не выполнился бы никогда.
Строки 12–15 (строки 40–49): WithCancellation(token). Токен пришёл не параметром метода, а при GetAsyncEnumerator(token). Он подставился в параметр с атрибутом [EnumeratorCancellation] (в нашем итераторе это ct). Без этого атрибута токен из WithCancellation пропал бы: параметр не получил бы его (предупреждение CS8425).
Строка 16 (строки 50–55): первый GetAsyncEnumerator вернул саму последовательность, как у обычного итератора.
Предскажите и задание (TODO 1)
В start/Program.cs у строк стоят комментарии PREDICT: сколько элементов придёт до отмены и какого типа будет исключение. Запустите и сверьтесь с таблицей выше:
cd chaptersdotnetrun-cRelease--project11-async-streams\start
Затем добавьте в CountAsync блок try/finally с печатью и проверьте, что finally выполняется и при отмене, и при досрочном break (добавьте второй цикл с break). А что будет, если не вызвать DisposeAsync руками, пользуясь GetAsyncEnumerator напрямую?
Строки 17–27 вывода (опыт «AsyncLocal», строки 57–64 кода): итератор записал Ambient = 99, увидел 99 перед первым yield return, но потребитель увидел 1, а перед вторым yield return внутри итератора снова 1. Запись не пережила даже один шаг. Дело в MoveNextAsync:
Сгенерированный MoveNextAsync и DisposeAsync (ILSpy, Release; фрагмент <CountAsync>d__1)
Строка 9: <>t__builder.MoveNext(ref stateMachine) — это AsyncMethodBuilderCore.Start (листинг ниже). Тот самый Start, который в главе 7, §7.3 возвращает вызывающему его ExecutionContext после синхронной части. Каждый MoveNextAsync — как вызов отдельного async-метода: что итератор записал в AsyncLocal, потребитель не увидит, и на следующем шаге итератор тоже не увидит своей записи.
Строки 10–14: если итератор дошёл до yield returnбез приостановки, статус уже Succeeded, и возвращается готовый ValueTask<bool> с результатом, без объекта (§11.3).
Строка 15: иначе в ValueTask кладётся сам итератор как источник (IValueTaskSource<bool>) с версией.
Строки 20–23:DisposeAsync в середине MoveNextAsync (состояние >= -1, то есть шаг идёт) бросает NotSupportedException: параллельно итерировать один перечислитель нельзя.
Строки 28–32:DisposeAsync включает режим закрытия (<>w__disposeMode) и пускает машину ещё раз: она выходит из тела через finally.
Факт: AsyncLocal, записанный в асинхронном итераторе, не переживает шаг
То же правило, что для обычного async-метода (глава 7, §7.3), но заметнее: там запись терялась на выходе из метода, а здесь теряется между соседними элементами. Если нужно протащить значение через итератор, передавайте его параметром или кладите в объект-держатель, а не в AsyncLocal. Читать AsyncLocal в итераторе можно: значение потребителя на каждом шаге видно.
11.3. Что генерирует компилятор для async + yield return¶
CountAsync после компиляции: класс и главные поля (ILSpy, Release)
Строка 1: он и последовательность (IAsyncEnumerable<int>), и перечислитель (IAsyncEnumerator<int>), и машина состояний (IAsyncStateMachine), и источник для ValueTask (IValueTaskSource<bool>). Один объект на обход.
Строки 4–5:AsyncIteratorMethodBuilder (builder для итераторов) и ManualResetValueTaskSourceCore<bool> — готовая структура из CoreLib, которая хранит результат MoveNextAsync и продолжение потребителя. Именно из неё выросла идея «ValueTask без Task».
Строка 6: текущий элемент.
Строка 7: режим закрытия для DisposeAsync.
Строка 8: если и параметр метода, и GetAsyncEnumerator(token) дали токены, GetAsyncEnumerator создаёт связанный источник (CreateLinkedTokenSource, глава 9, §9.5), а по концу итерации освобождает его.
Строки 10–11:<>3__ct — значение параметра из вызова метода, ct — рабочая копия для текущего перечислителя.
Состояния сгенерированной машины:
<>1__state
Значение
-2
создан, ещё не использован (можно отдать самого себя как перечислитель) или завершён
-3
перечислитель выдан, итерация не начата
-1
тело выполняется
0, 1, …
приостановлен на await номер N (как в главе 2)
-4
приостановлен на yield return: ждёт следующего MoveNextAsync
Строки 10–13:MoveNext — это AsyncMethodBuilderCore.Start (глава 7): отсюда терялись AsyncLocal в §11.2.
Строки 15–18: приостановка на await внутри итератора — тот же AsyncTaskMethodBuilder<VoidTaskResult>.AwaitUnsafeOnCompleted, что в главе 3: бокс (глава 1) создаётся на первой приостановке. Бокс — отдельный объект, но создаётся один раз на весь обход, а не на элемент.
Строки 20–29: конец итерации.
Сам yield return в MoveNext итератора выглядит так (фрагмент того же класса):
CountAsync.MoveNext: yield return и конец (фрагмент)
IL_00b0:awaiter.GetResult();<>2__current=<i>5__2;// 1. положить элемент в Currentnum=(<>1__state=-4);// 2. запомнить: ждём следующего MoveNextAsyncgotoIL_01a6;...IL_01a6:<>v__promiseOfValueOrEnd.SetResult(result:true);// 3. завершить «обещание» MoveNextAsync значением true// конец тела:<>1__state=-2;<>2__current=0;<>t__builder.Complete();<>v__promiseOfValueOrEnd.SetResult(result:false);// MoveNextAsync вернёт false
Строки 3–4 и 8:yield return — это запись в Current, смена состояния на -4 и SetResult(true) на ManualResetValueTaskSourceCore. Тот, кто ждёт MoveNextAsync, получает true.
Строки 11–14: конец тела или исключение (SetException) завершают то же «обещание» значением false или ошибкой.
Console.WriteLine();Console.WriteLine("== 3. Аллокации на элемент: 100 000 элементов ==");awaitMeasure("синхронный проход, yield return без await",async()=>{awaitforeach(int_inSources.SyncAsync(100_000)){}});awaitMeasure("с паузой await Task.Yield() на каждом элементе",async()=>{awaitforeach(int_inSources.YieldingAsync(100_000)){}});awaitMeasure("то же, ConfigureAwait(false)",async()=>{awaitforeach(int_inSources.YieldingAsync(100_000).ConfigureAwait(false)){}});
== 3. Аллокации на элемент: 100 000 элементов ==
синхронный проход, yield return без await : 0.00 байт на элемент
с паузой await Task.Yield() на каждом элементе : 0.00 байт на элемент
то же, ConfigureAwait(false) : -0.02 байт на элемент
Строка 2 вывода: итератор без await (все 100 000 MoveNextAsync завершаются синхронно): ноль байт на элемент. Из §11.2: готовый ValueTask<bool> без объекта.
Строка 3: итератор, который приостанавливается на каждом элементе, тоже ноль. Источник ValueTask — сам итератор (ManualResetValueTaskSourceCore, Reset() перед каждым шагом), а продолжение потребителя хранится в нём. Бокс для await Task.Yield() создаётся один раз (§11.3).
Строка 4: отрицательное число — шум измерения. Общая аллокация на 100 000 элементов — десятки байт на весь обход.
Сравните с вариантом, где каждый шаг возвращал бы Task<bool>: по объекту на элемент. Так выглядела бы цена асинхронного перечисления без ValueTask (глава 12).
Следствия для кода:
IAsyncEnumerable<T> дёшев на элемент. Но каждый обход создаёт объекты (итератор, возможно связанный источник токенов). Для «горячих» коротких последовательностей это заметно.
ConfigureAwait(false) для цикла:await foreach (var x in seq.ConfigureAwait(false)) применяется к MoveNextAsync и DisposeAsync потребителя (глава 5). Внутри самого итератора он нужен на его await, как в любом async-методе.
Не смешивайте с буферизацией.IAsyncEnumerable отдаёт по одному элементу и не читает вперёд: это стриминг из БД и HTTP (ASP.NET Core умеет отдавать IAsyncEnumerable<T> как поток JSON).
Раньше LINQ-операторы для IAsyncEnumerable<T> давал отдельный пакет System.Linq.Async. В .NET 10 они входят в платформу: сборка System.Linq.AsyncEnumerable лежит рядом с System.Linq в Microsoft.NETCore.App, и наша лабораторная использует Where, Select, ToListAsync без единого PackageReference.
== 4. LINQ над IAsyncEnumerable (BCL, .NET 10) ==
[итератор] finally
чётные ×10: 20, 40, 60
CountAsync(кратные 7 из 1000): 143
FirstAsync(> 500): 501
Where с async-предикатом: 0, 3, 6, 9
[итератор] finally
Take(3): 1, 2, 3 (источник остановлен: см. finally выше)
[итератор] finally
ToListAsync(ct) отменён
Методы: в сборке есть привычный набор (Where, Select, SelectMany, Take, Skip, Concat, Zip, OrderBy, GroupBy, Join, Distinct, Chunk, Reverse и другие) и «терминальные» операторы с суффиксом Async: ToListAsync, ToArrayAsync, ToDictionaryAsync, CountAsync, FirstAsync, AnyAsync, AggregateAsync, SumAsync, MinAsync/MaxAsync и т. д. Они возвращают ValueTask<T> и принимают CancellationToken.
Асинхронные предикаты: у большинства операторов есть перегрузка с Func<T, CancellationToken, ValueTask<bool>> (строка 80 кода: Where(async (n, ct) => …)), а не только с обычным Func<T, bool>.
Конец потребления освобождает источник:Take(3) остановился после третьего элемента, и в выводе [итератор] finally (строка 7) — DisposeAsync источника вызван. Так же ToListAsync(ct), отменённый токеном (строки 9–10): итератор закрылся.
Ленивость:Where и Select ничего не делают, пока не вызван терминальный оператор или await foreach.
Для перевода обычной последовательности в асинхронную: IEnumerable<T>.ToAsyncEnumerable().
Задание (TODO 3)
В start/Program.cs оставьте из CountAsync(10) только нечётные элементы и умножьте их на 10, одним выражением LINQ над IAsyncEnumerable<int> и с ToListAsync. Не подключайте никаких пакетов.
.NET 11
Сравнил открытые статические методы System.Linq.AsyncEnumerable отражением: в .NET 10.0.12 их 95, в .NET 11.0 RC1 (11.0.100-rc.1.26425.128) — 101. Разница только в соединениях: появился FullJoin (две перегрузки), а у Join, GroupJoin, LeftJoin, RightJoin — перегрузки без компаратора (в 10.0.12 только с компаратором, LeftJoin и RightJoin там уже есть). Остальной набор операторов тот же. Сравнение по числу параметров, без разбора типов; для релиза .NET 11 повторить. В этой книге все опыты сняты на рантайме 10.0.12.
11.6. Channel<T>: асинхронная очередь «производитель — потребитель»¶
Channel<T> — стандартный способ делать фоновую обработку в ASP.NET Core (очередь заданий + BackgroundService), замена BlockingCollection<T> в async-коде. Две стороны: Writer (TryWrite, WriteAsync, Complete) и Reader (TryRead, ReadAsync, ReadAllAsync, Completion). Какие каналы создаёт Channel:
System.Threading.Channels 10.0.12: фабрика Channel (сокращено)
Строки 1–4, 10–12, 19–21:все продолжения читателей по умолчанию асинхронные: runContinuationsAsynchronously = !AllowSynchronousContinuations, а AllowSynchronousContinuations по умолчанию false (§11.9). Это тот же принцип, что RunContinuationsAsynchronously у TaskCompletionSource (глава 5).
Строки 8–11:SingleReader = true выбирает другую реализацию, оптимизированную под одного читателя.
Строки 17–19: ёмкость 0 (и меньше) даёт RendezvousChannel: канал без буфера, запись завершается, только когда её принял читатель. В опыте 5в: capacity 0 (рандеву): RendezvousChannel1, TryWrite без читателя → False.
Строки 24–27:UnboundedPrioritizedChannel — канал с приоритетами (порядок по Comparer).
Типичный сценарий: производитель пишет, потребитель читает через await foreach:
varchannel=Channel.CreateBounded<int>(newBoundedChannelOptions(capacity:3){FullMode=BoundedChannelFullMode.Wait,// производитель будет ждать при переполненииSingleReader=true,});varproducer=Task.Run(async()=>{for(inti=1;i<=8;i++)awaitchannel.Writer.WriteAsync(i);// ждёт, если очередь полна (обратное давление)channel.Writer.Complete();// сигнал «больше не будет»});varconsumer=Task.Run(async()=>{awaitforeach(variteminchannel.Reader.ReadAllAsync())awaitTask.Delay(200);// медленный потребитель});awaitTask.WhenAll(producer,consumer);
ReadAllAsync — это обычный асинхронный итератор над теми же WaitToReadAsync и TryRead:
Строка 3: ждать, пока в канале что-то появится (или он закроется). ConfigureAwait(false) внутри.
Строки 5–8: выбрать из канала всё, что есть, синхронно и по одному отдавать потребителю. Пока в канале есть элементы, MoveNextAsync завершается синхронно, без приостановки (§11.4).
publicoverrideboolTryWrite(Titem){...lock(parent.SyncObj){if(parent._doneWriting!=null){returnfalse;}intcount=parent._items.Count;if(count!=0){if(count<parent._bufferedCapacity){parent._items.EnqueueTail(item);returntrue;}if(parent._mode==BoundedChannelFullMode.Wait){returnfalse;}if(parent._mode==BoundedChannelFullMode.DropWrite){parent._itemDropped?.Invoke(item);returntrue;}Tobj=((parent._mode==BoundedChannelFullMode.DropNewest)?parent._items.DequeueTail():parent._items.DequeueHead());parent._items.EnqueueTail(item);parent._itemDropped?.Invoke(obj);returntrue;}...// канал пуст: отдать элемент ждущему читателю или положить в очередь}}
Строки 6–9: после Complete() писать нельзя: false.
Строки 13–17: есть место — положить в хвост.
Строки 18–21: Wait — места нет, TryWrite возвращает false. WriteAsync в этом режиме вернёт ValueTask, который завершится, когда читатель освободит место (обратное давление: писатель замедляется до темпа читателя).
Строки 22–26: DropWrite — выбросить записываемый элемент и вернуть true (!). Вызывающий не узнает, что элемент потерян, кроме как через itemDropped.
Строки 27–30:DropNewest убирает последний в буфере (хвост), DropOldest — первый (голову), потом записывает новый.
== 5б. Режимы переполнения: 6 записей в канал на 3 элемента без читателя ==
DropOldest: TryWrite → TTTTTT, выброшено [1, 2, 3], осталось [4, 5, 6]
DropNewest: TryWrite → TTTTTT, выброшено [3, 4, 5], осталось [1, 2, 6]
DropWrite : TryWrite → TTTTTT, выброшено [4, 5, 6], осталось [1, 2, 3]
Wait : четвёртый TryWrite → False; WriteAsync вернул ValueTask, завершённый: False
Режим
Что выбрасывается при переполнении
Что остаётся из 1…6 при ёмкости 3
TryWrite
Wait
ничего: писатель ждёт
— (запись откладывается)
false при полном буфере
DropOldest
самый старый из буфера
4, 5, 6 (последние)
true
DropNewest
самый новый из буфера
1, 2, 6
true
DropWrite
пишущийся элемент
1, 2, 3 (первые)
true (!)
Факт: в режимах Drop*TryWrite всегда возвращает true
Даже если элемент тут же выброшен (DropWrite) или вытеснил другой. Единственный способ узнать о потере — колбэк itemDropped в Channel.CreateBounded(options, itemDropped). Колбэк вызывается вне блокировки канала, на потоке писателя: делайте его быстрым.
Задание (TODO 2)
В start/Program.cs замените FullMode = Wait на DropOldest. Что теперь прочитает потребитель? Запустите программу несколько раз: одинаков ли результат? Объясните разницу с помощью itemDropped.
(final/Program.cs, функция RunChannel; полный вывод — запуском.) Видно обратное давление: писатель быстро записал 1, 2, 3, после чего его запись 4 стала возможной, только когда читатель взял 1; дальше каждая запись приходит сразу после чтения (в темпе читателя, 200 мс). Все 8 элементов дошли.
С DropOldest писатель не ждёт никого: все 8 записей за миллисекунду, пять элементов вытеснены (1…5), читатель получает 6, 7, 8.
Факт: какие элементы переживут DropOldest, зависит от гонки
В 10 запусках сценария DropOldest читатель получил 6, 7, 8 в шести и 2, 6, 7, 8 в четырёх: в этих запусках он успел взять элемент 2 раньше, чем писатель вытеснил его. Это нормально: режим гарантирует «в канале всегда последние N», а не «читатель увидит именно эти». Если потребитель должен видеть конкретные элементы, DropOldest не подходит.
Строки 3–12: если токена отмены нет, используется один заранее созданный и повторно используемый объект_writerSingleton (TryOwnAndReset: «возьми, если свободен»). Поэтому ожидающая запись не аллоцирует.
Строки 13–16: если свободного синглтона нет или есть токен отмены, создаётся новый объект (pooled: false).
Строки 9 и 17: блокированный писатель — это очередь операций (_blockedWritersHead), а возвращаемое значение — ValueTask поверх объекта-источника. Читатель, освободив место, завершит операцию первого блокированного писателя.
== 5в. Завершение канала ==
после Complete: TryWrite → False; Completion завершена: False (остались непрочитанные)
WriteAsync в закрытый канал: ChannelClosedException
дочитано 1
дочитано 2
после дочитывания Completion: RanToCompletion
получено 1
Complete(exception): читатель получил «производитель упал» после дочитывания
capacity 0 (рандеву): RendezvousChannel`1, TryWrite без читателя → False
== 5г. Несколько читателей ==
30 заданий обработано читателями: 10 + 10 + 10 = 30; Completion: RanToCompletion
Строка 2 вывода:Complete() закрывает канал для записи (TryWrite → false), но не для чтения. Reader.Completion ещё не завершена, пока не прочитаны оставшиеся элементы.
Строка 3:WriteAsync в закрытый канал бросает ChannelClosedException.
Строки 4–6: читатель дочитывает остаток, и только потом Completion переходит в RanToCompletion. Так делается корректное завершение: производитель вызывает Complete(), потребитель выходит из await foreach сам.
Строки 7–8:Complete(exception) — так писатель сообщает об ошибке: читатель получает уже накопленные элементы и затем исключение из ReadAllAsync / ReadAsync.
Строка 12: три читателя поделили 30 заданий поровну (10 + 10 + 10), и все вышли из await foreach, когда писатель вызвал Complete(). Один канал на нескольких читателей работает как очередь заданий: каждое задание получает ровно один читатель (для этого SingleReader оставляют false).
11.9. AllowSynchronousContinuations: чужой код внутри вашего TryWrite¶
Когда читатель ждёт на пустом канале, а писатель делает TryWrite, нужно «разбудить» читателя. Как это делается, решает AsyncOperation.SignalCompletion:
privateprotectedvoidSignalCompletion(){...objectcapturedContext=_capturedContext;if((capturedContext==null||capturedContextisExecutionContext)?true:false){if(RunContinuationsAsynchronously){UnsafeQueueSetCompletionAndInvokeContinuation();// поставить в пулreturn;}}else{SynchronizationContextsynchronizationContext=...;if(synchronizationContext!=null){if(RunContinuationsAsynchronously||synchronizationContext!=SynchronizationContext.Current){synchronizationContext.Post(…);// вернуть в контекст читателяreturn;}}else{TaskSchedulertaskScheduler=...;if(RunContinuationsAsynchronously||taskScheduler!=TaskScheduler.Current){Task.Factory.StartNew(…,taskScheduler);// в планировщик читателяreturn;}}}SetCompletionAndInvokeContinuation();// иначе — выполнить продолжение прямо здесь}
Это та же логика выбора места продолжения, что в главе 5, только написанная заново для каналов:
Строки 5–11: читатель без контекста: продолжение ставится в пул, если асинхронные продолжения включены (по умолчанию).
Строки 15–23: читатель в SynchronizationContext — Post в него.
Строки 25–31: читатель в нестандартном TaskScheduler — задача в нём.
Строка 34: иначе (синхронные продолжения разрешены) — выполнить продолжение читателя прямо здесь, внутри вызова писателя.
Console.WriteLine();Console.WriteLine("== 6. AllowSynchronousContinuations ==");foreach(boolallowinnew[]{false,true}){varch=Channel.CreateUnbounded<int>(newUnboundedChannelOptions{SingleReader=true,AllowSynchronousContinuations=allow});Flag.InsideWrite=false;boolranInsideTryWrite=false;Taskreader=Task.Run(async()=>{awaitch.Reader.ReadAsync();// читатель ждёт на пустом каналеranInsideTryWrite=Flag.InsideWrite;Thread.Sleep(200);// «обработка» в продолжении читателя});awaitTask.Delay(100);longtryWriteMs=0;awaitTask.Run(()=>{varsw=Stopwatch.StartNew();Flag.InsideWrite=true;ch.Writer.TryWrite(1);Flag.InsideWrite=false;tryWriteMs=sw.ElapsedMilliseconds;});awaitreader;Console.WriteLine($" AllowSynchronousContinuations={allow,-5}: TryWrite занял {tryWriteMs,3} мс, продолжение читателя выполнилось внутри TryWrite: {ranInsideTryWrite}");}
== 6. AllowSynchronousContinuations ==
AllowSynchronousContinuations=False: TryWrite занял 0 мс, продолжение читателя выполнилось внутри TryWrite: False
AllowSynchronousContinuations=True : TryWrite занял 200 мс, продолжение читателя выполнилось внутри TryWrite: True
Читатель после ReadAsync «обрабатывает» 200 мс (строка 160 кода). С AllowSynchronousContinuations = true его обработка выполнилась внутри TryWrite писателя: TryWrite занял те же 200 мс, и писатель стоял всё это время (строка 3 вывода). По умолчанию (false) TryWrite вернулся сразу (строка 2), а читатель обработал элемент в пуле.
Включать AllowSynchronousContinuations можно, только если продолжение читателя гарантированно короткое и не блокирует (оптимизация: нет постановки в пул), и писатель не держит блокировок, которые продолжение может захотеть. В остальных случаях оставляйте false. Это то же правило, что с TaskCreationOptions.RunContinuationsAsynchronously (глава 5, §5.6) и SemaphoreSlim (глава 10, §10.6).
Console.WriteLine();Console.WriteLine("== 7. Стоимость элемента канала: 200 000 элементов, один писатель и один читатель ==");foreach(var(name,options)innew(string,BoundedChannelOptions)[]{("Bounded(64)",newBoundedChannelOptions(64)),("Bounded(64), SingleReader + SingleWriter",newBoundedChannelOptions(64){SingleReader=true,SingleWriter=true}),}){var(ms,bytes)=awaitPump(options);Console.WriteLine($" {name,-42}: {ms,4} мс, {bytes:F1} байт на элемент");}
== 7. Стоимость элемента канала: 200 000 элементов, один писатель и один читатель ==
Bounded(64) : 97 мс, 0.1 байт на элемент
Bounded(64), SingleReader + SingleWriter : 72 мс, 0.0 байт на элемент
Байты считаются по всему процессу (GC.GetTotalAllocatedBytes), то есть в обоих потоках. Около нуля байт на элемент: ожидающие операции канала — повторно используемые объекты (§11.7), ReadAllAsync — итератор без аллокаций на элемент (§11.4), значения int хранятся в Deque<T>. SingleReader и SingleWriter включают более лёгкие пути (меньше блокировок): в двух прогонах на 25–30 % быстрее (76 → 53 мс, 97 → 72 мс).
yield return и await — разные «паузы»: первая отдаёт значение вызывающему и возобновляется им, вторая ждёт уведомления. Асинхронный итератор — класс, объединяющий последовательность, перечислитель, машину состояний и источник для ValueTask<bool>.
await foreach всегда вызывает DisposeAsync: при break и исключениях выполняются finally итератора. Токен из WithCancellation попадает в параметр с [EnumeratorCancellation].
Каждый MoveNextAsync выполняется через AsyncMethodBuilderCore.Start: AsyncLocal, записанный внутри итератора, не переживает шаг.
Элемент асинхронного итератора не аллоцирует (ни синхронный, ни приостановленный): ValueTask поверх ManualResetValueTaskSourceCore.
В .NET 10 LINQ для IAsyncEnumerable входит в BCL (System.Linq.AsyncEnumerable), в том числе операторы с асинхронными предикатами.
// Глава 11, заготовка. Асинхронный итератор и ограниченный канал.// PREDICT: сколько элементов успеет прийти до отмены? Как будут чередоваться «записано» и «прочитано»?// TODO 1 (§11.2): добавьте в CountAsync блок try/finally с печатью и убедитесь, что finally выполняется при отмене и при break.// TODO 2 (§11.5): замените FullMode = Wait на DropOldest. Что прочитает потребитель? Одинаков ли результат от запуска к запуску?// TODO 3 (§11.4): отфильтруйте нечётные и умножьте на 10 с помощью LINQ над IAsyncEnumerable (без System.Linq.Async: он есть в BCL).usingSystem.Runtime.CompilerServices;usingSystem.Threading.Channels;Console.WriteLine("== IAsyncEnumerable ==");usingvarcts=newCancellationTokenSource(TimeSpan.FromMilliseconds(450));try{awaitforeach(varninCountAsync(10,cts.Token))Console.WriteLine($"получено {n}");// PREDICT: до какого n?}catch(OperationCanceledExceptione){Console.WriteLine($"отменено по таймауту: {e.GetType().Name}");// PREDICT: какой тип?}Console.WriteLine();Console.WriteLine("== Channel ==");varchannel=Channel.CreateBounded<int>(newBoundedChannelOptions(capacity:3){FullMode=BoundedChannelFullMode.Wait,// производитель будет ждать при переполненииSingleReader=true,});varproducer=Task.Run(async()=>{for(inti=1;i<=8;i++){awaitchannel.Writer.WriteAsync(i);Console.WriteLine($" записано {i}");}channel.Writer.Complete();});varconsumer=Task.Run(async()=>{awaitforeach(variteminchannel.Reader.ReadAllAsync()){Console.WriteLine($"прочитано {item}");awaitTask.Delay(200);// медленный потребитель}});awaitTask.WhenAll(producer,consumer);staticasyncIAsyncEnumerable<int>CountAsync(intmax,[EnumeratorCancellation]CancellationTokenct=default){for(inti=1;i<=max;i++){awaitTask.Delay(100,ct);yieldreturni;}}
usingSystem.Runtime.CompilerServices;// Итераторы для опытов главы 11.staticclassSources{// Обычный (синхронный) итератор: yield return без async.publicstaticIEnumerable<int>Numbers(intmax){Console.WriteLine(" [итератор] тело началось");for(inti=1;i<=max;i++)yieldreturni;Console.WriteLine(" [итератор] тело закончилось");}// Асинхронный итератор с настоящей паузой перед каждым элементом.publicstaticasyncIAsyncEnumerable<int>CountAsync(intmax,[EnumeratorCancellation]CancellationTokenct=default){try{for(inti=1;i<=max;i++){awaitTask.Delay(100,ct);yieldreturni;}}finally{Console.WriteLine(" [итератор] finally");}}// Асинхронный итератор, который ни разу не приостанавливается: yield return без await.publicstaticasyncIAsyncEnumerable<int>SyncAsync(intmax){for(inti=0;i<max;i++)yieldreturni;awaitTask.CompletedTask;}// Асинхронный итератор, который приостанавливается на каждом элементе.publicstaticasyncIAsyncEnumerable<int>YieldingAsync(intmax){for(inti=0;i<max;i++){awaitTask.Yield();yieldreturni;}}publicstaticreadonlyAsyncLocal<int>Ambient=new();// Итератор записывает AsyncLocal и смотрит на него перед каждым yield return.publicstaticasyncIAsyncEnumerable<int>WritesAmbientAsync(boolpause){Console.WriteLine($" [итератор] до записи Ambient = {Ambient.Value}");Ambient.Value=99;for(inti=0;i<2;i++){if(pause)awaitTask.Yield();Console.WriteLine($" [итератор] перед yield Ambient = {Ambient.Value}");yieldreturni;}}}