Rx.NET中如何实现类TaskCompletionSource行为的Subject
Rx.NET 如何实现与TaskCompletionSource行为一致的Subject
需求说明
需要实现一个具备如下行为的自定义Subject(命名为PromiseSubject):
- 仅缓存序列产生的第一个
OnNext事件值 - 无论序列后续是否调用
OnCompleted终止,所有订阅者(包括序列终止后才发起订阅的观察者)都能收到这个缓存的第一个值 - 第一个值之后产生的所有
OnNext事件直接忽略,不做缓存也不向下游推送
内置实现无法满足需求:
AsyncSubject只会在序列完成后才推送最终值,缓冲区大小为1的ReplaySubject(bufferSize:1)会缓存最后一个值而非第一个值,二者行为均不符合预期。
预期行为示例
// 对比内置ReplaySubject的行为差异 //var subj = new ReplaySubject<int>(bufferSize: 1); var subj = new PromiseSubject<int>(); // 序列完成前订阅 subj.Subscribe(i => Console.WriteLine(i)); subj.OnNext(1); subj.OnNext(2); subj.OnNext(3); subj.OnCompleted(); // 序列完成后订阅 subj.Subscribe(i => Console.WriteLine(i)); Console.ReadLine();
预期输出结果:
1 1
现有实现参考
目前基于TaskCompletionSource封装的可运行实现如下:
public class PromiseSubject<T> : ISubject<T> { private readonly TaskCompletionSource<(bool HasValue, T Value)> _tcs; private readonly IObservable<T> _observable; public PromiseSubject() { _tcs = new(TaskCreationOptions.RunContinuationsAsynchronously); _observable = _tcs.Task.ToObservable() .Where(r => r.HasValue).Select(r => r.Value!); } public void OnCompleted() => _tcs.TrySetResult((false, default!)); public void OnError(Exception error) => _tcs.TrySetException(error); public void OnNext(T value) => _tcs.TrySetResult((true, value)); public IDisposable Subscribe(IObserver<T> observer) => _observable.Subscribe(observer); }
核心疑问
是否存在更优雅的实现方式,不需要自定义Subject类封装TaskCompletionSource,直接通过组合Rx.NET内置运算符即可实现该行为?
内容的提问来源于stack exchange,提问作者noseratio
相关产品推荐
相关产品推荐

