Енумераторы в дотнете можно использовать двумя способами: через foreach, и вручную, через MoveNext/Dispose. Второй способ часто приводит к багам.
В этом методе разработчик пытался реализовать обёртку для
IAsyncEnumerable<T>, которая распараллеливает обработку текущего элемента последовательности и получение следующего:public static async IAsyncEnumerable<T> Prefetch<T>(
this IAsyncEnumerable<T> source) {
await using var enumerator = source.GetAsyncEnumerator();
var nextTask = enumerator.MoveNextAsync();
while (await nextTask) {
var current = enumerator.Current;
nextTask = enumerator.MoveNextAsync(); // запускаем MoveNext перед возвратом управления в вызывающий код
yield return current;
}
}
// USAGE //
var sw = Stopwatch.StartNew();
await foreach (var x in En())
await Task.Delay(100);
Console.WriteLine("Without prefetch: " + sw.ElapsedMilliseconds); // 700
sw.Restart();
await foreach (var x in En().Prefetch())
await Task.Delay(100);
Console.WriteLine("With prefetch: " + sw.ElapsedMilliseconds); // 400
async IAsyncEnumerable<int> En()
{
await Task.Delay(100);
yield return 1;
await Task.Delay(100);
yield return 2;
await Task.Delay(100);
yield return 3;
await Task.Delay(100);
}
Однако, такой код сломается, если внутри
await foreach сделать break:await foreach (var x in En().Prefetch())
break;
Unhandled exception. System.NotSupportedException: Specified method is not supported.
at Program.<<Main>$>g__En|0_0()+System.IAsyncDisposable.DisposeAsync()
at Extensions.Prefetch[T](IAsyncEnumerable`1 source)+MoveNext()
Ошибка возникает потому что у `IAsyncEnumerator<T>` нельзя одновременно вызывать `MoveNextAsync` и `DisposeAsync`. А в реализации
Prefetch так и происходит — после break сразу вызывается DisposeAsync (из-за await using) для переменной enumerator, при наличии незавершенного nextTaskВ зависимости от реализации
IAsyncEnumerator<T> возможны и другие симптомы неправильного взаимодействия MoveNext и Dispose, например утечки памятиИсправить можно явно дожидаясь
nextTask перед enumerator.DisposeAsync(). При этом нужно учесть, что ValueTask можно await-ить только один раз:public static async IAsyncEnumerable<T> Prefetch<T>(this IAsyncEnumerable<T> source)
{
await using var enumerator = source.GetAsyncEnumerator();
var nextTask = enumerator.MoveNextAsync();
try {
while (true) {
var nextTaskResult = await nextTask;
nextTask = default; // предотвращаем повторный await ValueTask
if (!nextTaskResult)
yield break;
var current = enumerator.Current;
nextTask = enumerator.MoveNextAsync();
yield return current;
}
}
finally {
if (nextTask != default)
await nextTask;
}
}
Такое исправление решает эту проблему — теперь
MoveNextAsync и DisposeAsync вызываются последовательно, и ошибки не возникает.@epeshkblog