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());
核心逻辑说明
- 立即执行规则:当检测到当前不在冷却状态时,立即执行
Save()并启动2秒冷却计时器。 - 合并请求规则:冷却期间收到的所有请求都会被标记为“待处理”,冷却结束后自动执行一次
Save(),并重新进入冷却状态(若后续仍有新请求)。
内容的提问来源于stack exchange,提问作者user007
相关产品推荐
相关产品推荐

