【问题标题】:Processing a Stream into IAsyncEnumerable - Stream is not readable将流处理为 IAsyncEnumerable - 流不可读
【发布时间】:2021-03-11 09:24:59
【问题描述】:

我有一个工作流程,我尝试执行以下操作:

  • 一个接受回调的方法,它在内部产生一个Stream,并且该方法的调用者可以使用回调以任何他们想要的方式处理Stream
  • 在一种特殊情况下,调用者使用回调从 Stream 中生成 IAsyncEnumerable

我在下面创建了一个最小的复制示例:

class Program
{
    private static async Task<Stream> GetStream()
    {
        var text =
            @"Multi-line
            string";

        await Task.Yield();

        var bytes = Encoding.UTF8.GetBytes(text);
        return new MemoryStream(bytes);
    }

    private static async Task<T> StreamData<T>(Func<Stream, T> streamAction)
    {
        await using var stream = await GetStream();
        return streamAction(stream);
    }

    private static async Task StreamData(Func<Stream, Task> streamAction)
    {
        await using var stream = await GetStream();
        await streamAction(stream);
    }

    private static async IAsyncEnumerable<string> GetTextLinesFromStream(Stream stream)
    {
        using var reader = new StreamReader(stream);

        var line = await reader.ReadLineAsync();
        while (line != null)
        {
            yield return line;
            line = await reader.ReadLineAsync();
        }
    }

    private static async Task Test1()
    {
        async Task GetRecords(Stream str)
        {
            await foreach(var line in GetTextLinesFromStream(str))
                Console.WriteLine(line);
        }

        await StreamData(GetRecords);
    }

    private static async Task Test2()
    {
        await foreach(var line in await StreamData(GetTextLinesFromStream))
            Console.WriteLine(line);
    }

    static async Task Main(string[] args)
    {
        await Test1();
        await Test2();
    }
}  

在这里,方法Test1 可以正常工作,而Test2 不能,因为Stream is not readable 失败。问题是在第二种情况下,当代码开始处理实际流时,流已经被释放了。

大概这两个例子的区别在于,对于第一个例子,读取流是在一次性stream的上下文中执行的,而在第二个例子中,我们已经退出了。

但是,我认为第二种情况也可能有效 - 至少我觉得它非常符合 C# 习惯。为了让第二个案例也能正常工作,我还有什么遗漏吗?

【问题讨论】:

    标签: c# .net asynchronous iasyncenumerable


    【解决方案1】:

    Test2 方法的问题在于 Stream 在创建 IAsyncEnumerable&lt;string&gt; 时被释放,而不是在其枚举完成时释放。

    Test2 方法使用第一个 StreamData 重载,即返回 Task&lt;T&gt; 的重载。在这种情况下,TIAsyncEnumerable&lt;string&gt;。所以StreamData 方法返回一个产生异步序列的任务,然后立即释放流(在产生序列之后)。显然,这不是处理流的正确时机。正确的时机应该是在 await foreach 循环完成之后。

    为了使Test2 透明地工作,您应该添加返回Task&lt;IAsyncEnumerable&lt;T&gt;&gt;StreamData 方法的第三个重载(而不是TaskTask&lt;T&gt;)。此重载应返回与可处置资源绑定的专用异步序列,并在其枚举完成时处置此资源。下面是这样一个序列的实现:

    public class AsyncEnumerableDisposable<T> : IAsyncEnumerable<T>
    {
        private readonly IAsyncEnumerable<T> _source;
        private readonly IAsyncDisposable _disposable;
    
        public AsyncEnumerableDisposable(IAsyncEnumerable<T> source,
            IAsyncDisposable disposable)
        {
            // Arguments validation omitted
            _source = source;
            _disposable = disposable;
        }
    
        async IAsyncEnumerator<T> IAsyncEnumerable<T>.GetAsyncEnumerator(
            CancellationToken cancellationToken)
        {
            await using (_disposable.ConfigureAwait(false))
                await foreach (var item in _source
                    .WithCancellation(cancellationToken)
                    .ConfigureAwait(false)) yield return item;
        }
    }
    

    您可以像这样在StreamData 方法中使用它:

    private static async Task<IAsyncEnumerable<T>> StreamData<T>(
        Func<Stream, IAsyncEnumerable<T>> streamAction)
    {
        var stream = await GetStream();
        return new AsyncEnumerableDisposable<T>(streamAction(stream), stream);
    }
    

    请记住,通常IAsyncEnumerable&lt;T&gt; 可以在其生命周期内多次枚举,并且通过将其包装到AsyncEnumerableDisposable&lt;T&gt; 中,它基本上被简化为单枚举序列(因为资源将在第一个枚举)。


    替代方案:System.Interactive.Async 包包含 AsyncEnumerableEx.Using 运算符,可用于代替自定义 AsyncEnumerableDisposable 类:

    private static async Task<IAsyncEnumerable<T>> StreamData<T>(
        Func<Stream, IAsyncEnumerable<T>> streamAction)
    {
        var stream = await GetStream();
        return AsyncEnumerableEx.Using(() => stream, streamAction);
    }
    

    不同之处在于Stream 将通过其Dispose 方法同步释放。 AFAICS 不支持在此包中处理 IAsyncDisposables。

    这是AsyncEnumerableEx.Using方法的签名:

    // Constructs an async-enumerable sequence that depends on a resource object, whose
    // lifetime is tied to the resulting async-enumerable sequence's lifetime.
    public static IAsyncEnumerable<TSource> Using<TSource, TResource>(
        Func<TResource> resourceFactory,
        Func<TResource, IAsyncEnumerable<TSource>> enumerableFactory)
        where TResource : IDisposable;
    

    【讨论】:

    • 感谢您的详细解释和建议的解决方案 - 我倾向于同意这可能是最接近我想要实现的目标。太糟糕了,它需要对 StreamData 方法进行更专业的覆盖。无论如何,接受这个作为答案。
    • @zidour 是的,我不认为你可以只用Task&lt;T&gt; StreamData&lt;T&gt;(Func&lt;Stream, T&gt; streamAction) 超载。这太笼统了,并且没有提供一种机制来指示处理 Stream 的正确时间。
    猜你喜欢
    • 2011-04-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2010-11-07
    • 2011-07-20
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多