【问题标题】:Implementation of `Stream` that sends to IAsyncEnumerable<bytes>发送到 IAsyncEnumerable<bytes> 的 `Stream` 的实现
【发布时间】:2021-07-05 11:34:03
【问题描述】:

在我的 .Net Core 应用程序中,第 3 方库中有一个方法写入 System.IO.Stream 接口(它将流接口作为参数并写入它),但我希望该数据进入我的数据期望数据为IAsyncEnumerable&lt;bytes&gt; 流的源。我着手编写实现Stream 接口的代码,所以当Write() 被调用时,它会写入IAsyncEnumerable&lt;bytes&gt;,然后认为“这一定是之前完成的” - 似乎它是通用的。

那么在 3rd 方库中是否有标准实现,或者我缺少任何“巧妙的技巧”?

【问题讨论】:

  • 有一个Stream.ReadAsync 方法可以用来解决这个问题,但不是所有的重载都支持旧的.NET 平台。您的目标是什么 .NET 平台?写入流的byte 是否必须立即从IAsyncEnumerable&lt;byte&gt; 浮出水面,或者您可以通过一些延迟(缓冲)来提高性能?
  • 谢谢,但Stream.ReadAsync 将用于读取流,而我的问题是某些内容正在写入流,我想放入它​​正在写入的流的实现,然后写入我的应用程序需要 IAsyncEnumerable&lt;bytes&gt; 的字节。在现代 .Net Core 中开发(我在原始描述中错过了!)并且不介意延迟。我可以编写一个支持WriteStream 实现,我只是问它是否已经完成,因为我认为“不同类型的流之间的转换”将是一个常见的要求。
  • 这个Stream的出处是什么?第 3 方库是直接公开 Stream,还是允许您将自己的 Stream 实现作为参数传递给方法?在第二种情况下,您可以查看this 问题,作为解决此问题的第一步。
  • 感谢您的提问,Theodor - 这是第二种情况(我已经在问题中澄清了这一点)。链接的问题很有趣,它有一些很好的代码片段。我特别问是否已经有一个标准的实现,但我想没有,所以我必须自己写。谢谢

标签: .net .net-core iasyncenumerable


【解决方案1】:

这是一个自定义的Stream 实现,用于异步生产者-消费者场景。这是一个只可写的流,只有通过特殊的GetConsumingEnumerable 方法才能读取(消费)它。

public class ProducerConsumerStream : Stream
{
    private readonly Channel<byte> _channel;

    public ProducerConsumerStream(bool singleReader = true, bool singleWriter = true)
    {
        _channel = Channel.CreateUnbounded<byte>(new UnboundedChannelOptions()
        {
            SingleReader = singleReader,
            SingleWriter = singleWriter
        });
    }

    public override bool CanRead { get { return false; } }
    public override bool CanSeek { get { return false; } }
    public override bool CanWrite { get { return true; } }
    public override long Length { get { throw new NotSupportedException(); } }
    public override void Flush() { }

    public override long Position
    {
        get { throw new NotSupportedException(); }
        set { throw new NotSupportedException(); }
    }

    public override long Seek(long offset, SeekOrigin origin)
        => throw new NotSupportedException();

    public override void SetLength(long value)
        => throw new NotSupportedException();

    public override int Read(byte[] buffer, int offset, int count)
        => throw new NotSupportedException();

    public override void Write(byte[] buffer, int offset, int count)
    {
        if (buffer == null) throw new ArgumentNullException(nameof(buffer));
        if (offset < 0) throw new ArgumentOutOfRangeException(nameof(offset));
        if (count < 0) throw new ArgumentOutOfRangeException(nameof(count));
        if (offset + count > buffer.Length)
            throw new ArgumentOutOfRangeException(nameof(count));

        for (int i = offset; i < offset + count; i++)
            _channel.Writer.TryWrite(buffer[i]);
    }

    public override void WriteByte(byte value)
    {
        _channel.Writer.TryWrite(value);
    }

    public override void Close()
    {
        base.Close();
        _channel.Writer.Complete();
    }

    public IAsyncEnumerable<byte> GetConsumingEnumerable(
        CancellationToken cancellationToken = default)
    {
        return _channel.Reader.ReadAllAsync(cancellationToken);
    }
}

此实现基于Channel&lt;byte&gt;。如果对频道不熟悉,有教程here

【讨论】:

  • 谢谢,那我得根据这样的片段自己写了,因为我想这不是标准要求。
猜你喜欢
  • 2020-06-02
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-04-07
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多