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

如何配置C# Reactive ReplaySubject,使缓存项被访问后清除

实现一次性重播的自定义Subject(满足你的需求)

嘿,这个需求原生的ReplaySubject<T>确实做不到——它的缓存是持久化的,只要没触发清理条件(比如大小/时间限制),后续订阅者还是能拿到之前缓存的项。不过我们可以通过自定义一个ISubject<T>来实现你要的行为:无订阅时缓存项,有订阅者时一次性推送所有缓存并清空,后续订阅者拿不到已推送过的项。

自定义Subject实现

下面是一个线程安全的实现,核心是用队列缓存无订阅时的项,新订阅者进来时先消费所有缓存再订阅内部Subject:

using System;
using System.Collections.Generic;
using System.Reactive.Subjects;
using System.Reactive;

public class OneTimeReplaySubject<T> : ISubject<T>
{
    private readonly Subject<T> _innerSubject = new Subject<T>();
    private readonly Queue<T> _messageCache = new Queue<T>();
    private readonly object _lock = new object();

    public void OnNext(T value)
    {
        lock (_lock)
        {
            // 如果有活跃订阅者,直接发送;否则缓存
            if (_innerSubject.HasObservers)
            {
                _innerSubject.OnNext(value);
            }
            else
            {
                _messageCache.Enqueue(value);
            }
        }
    }

    public void OnError(Exception error)
    {
        lock (_lock)
        {
            _messageCache.Clear();
            _innerSubject.OnError(error);
        }
    }

    public void OnCompleted()
    {
        lock (_lock)
        {
            // 完成前先推送所有缓存项
            foreach (var item in _messageCache)
            {
                _innerSubject.OnNext(item);
            }
            _messageCache.Clear();
            _innerSubject.OnCompleted();
        }
    }

    public IDisposable Subscribe(IObserver<T> observer)
    {
        lock (_lock)
        {
            // 新订阅者先接收所有缓存项
            foreach (var cachedItem in _messageCache)
            {
                observer.OnNext(cachedItem);
            }
            // 清空缓存,避免后续订阅者拿到这些项
            _messageCache.Clear();
            
            // 订阅内部Subject,接收后续的实时消息
            return _innerSubject.Subscribe(observer);
        }
    }
}

使用示例

我们来验证一下这个Subject的行为是否符合你的需求:

var subject = new OneTimeReplaySubject<int>();

// 阶段1:无订阅者,发送的项被缓存
subject.OnNext(1);
subject.OnNext(2);
subject.OnNext(3);

// 阶段2:第一个订阅者进来,收到所有缓存项
var subscriber1 = subject.Subscribe(
    x => Console.WriteLine($"订阅者1收到:{x}"),
    ex => Console.WriteLine($"订阅者1出错:{ex.Message}"),
    () => Console.WriteLine("订阅者1完成")
);
// 输出:
// 订阅者1收到:1
// 订阅者1收到:2
// 订阅者1收到:3

// 阶段3:有活跃订阅者,新发送的项直接推送
subject.OnNext(4);
// 输出:订阅者1收到:4

// 阶段4:取消订阅,再次无订阅者,发送的项被缓存
subscriber1.Dispose();
subject.OnNext(5);
subject.OnNext(6);

// 阶段5:第二个订阅者进来,收到最新的缓存项(1-3已经被清空了)
var subscriber2 = subject.Subscribe(
    x => Console.WriteLine($"订阅者2收到:{x}"),
    () => Console.WriteLine("订阅者2完成")
);
// 输出:
// 订阅者2收到:5
// 订阅者2收到:6

// 阶段6:完成Subject,推送剩余缓存(这里没有)并通知完成
subject.OnCompleted();
// 输出:订阅者2完成

subscriber2.Dispose();

关键行为说明

  • 无订阅时缓存:当没有任何订阅者时,所有OnNext的项都会被存入内部队列。
  • 订阅时一次性推送并清空:新订阅者会立即收到所有缓存项,之后缓存被清空,后续订阅者无法获取这些已推送的项。
  • 线程安全:用lock保证了多线程环境下的缓存操作和订阅操作的原子性,符合Rx的线程安全契约。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 11:52:08