如何基于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
相关产品推荐
相关产品推荐

