You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何基于IObservable<byte>创建Stream?

如何在IObservable与Stream之间搭建高效桥梁

我一直在思考如何搭建IObservable<byte>与Stream之间的桥梁,但毫无头绪。使用场景是:某个函数要求传入Stream,但数据却存储在IObservable<byte>中。

由于Stream是拉取式(pull-based)、IObservable是推送式(push-based),我甚至不确定不调用byteSequence.ToEnumerable()是否可行。目前我能想到的实现方式是:

new MemoryStream(byteSequence.ToEnumerable().ToArray());

这种方法虽然可行,但会把所有字节都分配到内存中,在很多场景下并不适用。


更新:基于System.IO.Pipelines的实现

受Marc Gravell的评论启发,我尝试了用Pipes实现的代码,目前看起来能正常运行,但不确定是否存在问题,现将代码贴出供参考:

public static Stream ToStream(this IObservable<byte> observable, int bufferSize = 4096)
{
    var pipe = new System.IO.Pipelines.Pipe();

    var reader = pipe.Reader;
    var writer = pipe.Writer;

    observable
        .Buffer(bufferSize)
        .Select(buffer => Observable.FromAsync(async () => await writer.WriteAsync(buffer.ToArray())))
        .Concat()
        .Subscribe(_ => { }, onCompleted: () => { writer.Complete(); }, onError: exception => writer.Complete(exception));

    return reader.AsStream();
}

代码优化建议

  • 避免不必要的数组分配:buffer.ToArray()会每次创建新数组,可直接使用new ReadOnlyMemory<byte>(buffer)复用Buffer返回的列表内存,减少内存开销
  • 处理订阅生命周期:保留订阅对象,在Stream释放时取消订阅,防止内存泄漏
  • 优化异步写入流程:简化Select+Concat的写法,同时增加FlushAsync确保数据及时写入Pipe
  • 配置Pipe参数:根据场景调整PipeOptions的暂停/恢复阈值,平衡内存占用和写入效率

优化后的代码示例:

public static Stream ToStream(this IObservable<byte> observable, int bufferSize = 4096)
{
    var pipeOptions = new PipeOptions(
        pauseWriterThreshold: bufferSize * 2,
        resumeWriterThreshold: bufferSize
    );
    var pipe = new System.IO.Pipelines.Pipe(pipeOptions);
    var reader = pipe.Reader;
    var writer = pipe.Writer;

    var subscription = observable
        .Buffer(bufferSize)
        .SelectMany(buffer => Observable.FromAsync(async () =>
        {
            await writer.WriteAsync(new ReadOnlyMemory<byte>(buffer));
            await writer.FlushAsync();
        }))
        .Subscribe(
            _ => { },
            exception => writer.Complete(exception),
            () => writer.Complete()
        );

    return new PipeStreamWrapper(reader, writer, subscription);
}

// 自定义Stream包装类,处理资源释放逻辑
private class PipeStreamWrapper : Stream
{
    private readonly Stream _innerStream;
    private readonly PipeWriter _writer;
    private readonly IDisposable _subscription;

    public PipeStreamWrapper(PipeReader reader, PipeWriter writer, IDisposable subscription)
    {
        _innerStream = reader.AsStream();
        _writer = writer;
        _subscription = subscription;
    }

    public override bool CanRead => _innerStream.CanRead;
    public override bool CanSeek => _innerStream.CanSeek;
    public override bool CanWrite => _innerStream.CanWrite;
    public override long Length => _innerStream.Length;
    public override long Position { get => _innerStream.Position; set => _innerStream.Position = value; }

    public override void Flush() => _innerStream.Flush();
    public override int Read(byte[] buffer, int offset, int count) => _innerStream.Read(buffer, offset, count);
    public override long Seek(long offset, SeekOrigin origin) => _innerStream.Seek(offset, origin);
    public override void SetLength(long value) => _innerStream.SetLength(value);
    public override void Write(byte[] buffer, int offset, int count) => _innerStream.Write(buffer, offset, count);

    protected override void Dispose(bool disposing)
    {
        if (disposing)
        {
            _subscription.Dispose();
            _writer.Complete();
            _innerStream.Dispose();
        }
        base.Dispose(disposing);
    }
}

内容的提问来源于stack exchange,提问作者SuperJMN

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.07 01:42:27