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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.02 22:16:02