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

Rx中Save()操作实现:按2秒规则立即/延迟单次执行

解决方案

你需要的是**“立即执行冷却后的第一个请求,冷却期间的所有请求合并为一次,在冷却结束时执行”**的逻辑,单纯的Throttle无法满足需求——它需要等待事件完全停止指定时长才会触发,而你需要即使事件持续到来,也能在冷却周期结束时执行一次保存。以下是两种可行的实现方式:

方式一:直观状态跟踪(易理解)

通过外部变量跟踪冷却状态与待处理请求,逻辑清晰,适合快速落地:

var i = 0;
var saveSource = Observable.Interval(TimeSpan.FromMilliseconds(200)).Select(x => i++);

var coolDownPeriod = TimeSpan.FromSeconds(2);
bool isCoolingDown = false;
bool hasPendingSave = false;
object lockObj = new object(); // 保证线程安全

saveSource.Subscribe(_ =>
{
    lock (lockObj)
    {
        if (!isCoolingDown)
        {
            // 距离上次执行超过2秒,立即执行Save
            Save();
            isCoolingDown = true;
            
            // 启动冷却计时器,2秒后检查是否有待处理请求
            Observable.Timer(coolDownPeriod)
                .ObserveOn(Scheduler.CurrentThread) // 根据业务场景切换调度器
                .Subscribe(__ =>
                {
                    lock (lockObj)
                    {
                        isCoolingDown = false;
                        if (hasPendingSave)
                        {
                            // 冷却期间有请求,执行一次Save并重新进入冷却
                            Save();
                            hasPendingSave = false;
                            isCoolingDown = true;
                            
                            Observable.Timer(coolDownPeriod)
                                .ObserveOn(Scheduler.CurrentThread)
                                .Subscribe(___ =>
                                {
                                    lock (lockObj)
                                    {
                                        isCoolingDown = false;
                                    }
                                });
                        }
                    }
                });
        }
        else
        {
            // 冷却期间收到请求,标记为待处理
            hasPendingSave = true;
        }
    }
});

方式二:纯Rx状态管理(无外部变量)

使用Rx的Scan操作符跟踪状态,避免外部变量的线程安全问题,更符合Rx的响应式编程思想:

var i = 0;
var saveSource = Observable.Interval(TimeSpan.FromMilliseconds(200)).Select(x => i++);

var coolDownPeriod = TimeSpan.FromSeconds(2);
// 状态对象:是否处于冷却中、是否有待处理请求
var initialState = (IsCoolingDown: false, HasPending: false);

var saveTrigger = saveSource
    .Scan(initialState, (state, _) =>
    {
        if (!state.IsCoolingDown)
        {
            // 冷却完成,触发立即执行,进入冷却状态
            return (IsCoolingDown: true, HasPending: false);
        }
        else
        {
            // 冷却中,标记待处理请求
            return (IsCoolingDown: true, HasPending: true);
        }
    })
    .Publish(sharedState =>
        // 触发立即执行的信号:从非冷却进入冷却状态且无待处理
        sharedState.Where(s => s.IsCoolingDown && !s.HasPending)
            .Select(_ => Unit.Default)
            // 触发冷却结束后的执行信号:冷却结束且有待处理请求
            .Merge(
                sharedState
                    .Where(s => s.IsCoolingDown && !s.HasPending)
                    .SelectMany(_ => Observable.Timer(coolDownPeriod))
                    .WithLatestFrom(sharedState, (_, s) => s)
                    .Where(s => s.HasPending)
                    .Select(_ => Unit.Default)
            )
    )
    .DistinctUntilChanged();

saveTrigger.Subscribe(_ => Save());

核心逻辑说明

  1. 立即执行规则:当检测到当前不在冷却状态时,立即执行Save()并启动2秒冷却计时器。
  2. 合并请求规则:冷却期间收到的所有请求都会被标记为“待处理”,冷却结束后自动执行一次Save(),并重新进入冷却状态(若后续仍有新请求)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 22:53:10