IAsyncEnumerator и его использование вручную
Енумераторы в дотнете можно использовать двумя способами: через foreach, и вручную, через MoveNext/Dispose. Второй способ часто приводит к багам.
В этом методе разработчик пытался реализовать обёртку для IAsyncEnumerable, которая распараллеливает обработку текущего элемента последовательности и получение следующего:
public static async IAsyncEnumerable Prefetch(
this IAsyncEnumerable 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 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.g__En|0_0()+System.IAsyncDisposable.DisposeAsync()
at Extensions.Prefetch[T](IAsyncEnumerable`1 source)+MoveNext()
Ошибка возникает потому что у `IAsyncEnumerator` нельзя одновременно вызывать `MoveNextAsync` и `DisposeAsync`. А в реализации Prefetch так и происходит — после break сразу вызывается DisposeAsync (из-за await using) для переменной enumerator, при наличии незавершенного nextTask
В зависимости от реализации IAsyncEnumerator возможны и другие симптомы неправильного взаимодействия MoveNext и Dispose, например утечки памяти
Исправить можно явно дожидаясь nextTask перед enumerator.DisposeAsync(). При этом нужно учесть, что ValueTask можно await-ить только один раз:
public static async IAsyncEnumerable Prefetch(this IAsyncEnumerable 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
Енумераторы в дотнете можно использовать двумя способами: через foreach, и вручную, через MoveNext/Dispose. Второй способ часто приводит к багам.
В этом методе разработчик пытался реализовать обёртку для IAsyncEnumerable, которая распараллеливает обработку текущего элемента последовательности и получение следующего:
public static async IAsyncEnumerable Prefetch(
this IAsyncEnumerable 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 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.g__En|0_0()+System.IAsyncDisposable.DisposeAsync()
at Extensions.Prefetch[T](IAsyncEnumerable`1 source)+MoveNext()
Ошибка возникает потому что у `IAsyncEnumerator` нельзя одновременно вызывать `MoveNextAsync` и `DisposeAsync`. А в реализации Prefetch так и происходит — после break сразу вызывается DisposeAsync (из-за await using) для переменной enumerator, при наличии незавершенного nextTask
В зависимости от реализации IAsyncEnumerator возможны и другие симптомы неправильного взаимодействия MoveNext и Dispose, например утечки памяти
Исправить можно явно дожидаясь nextTask перед enumerator.DisposeAsync(). При этом нужно учесть, что ValueTask можно await-ить только один раз:
public static async IAsyncEnumerable Prefetch(this IAsyncEnumerable 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