如何实现过滤瞬时特殊事件、仅推送持续指定时长同类事件的Rx.NET运算符
Rx.NET 特殊事件持久化判定运算符实现
预期行为
- 普通类型事件直接向下游透传
- 特殊类型事件需观测到持续指定时长后再向下游推送,时长从第一个收到的特殊事件时间开始计算,达到阈值后所有后续特殊事件均直接透传,无需再次判定;若阈值时间内出现普通事件,自动重置特殊事件判定状态
通用扩展方法实现
直接封装为IObservable<T>的扩展方法,支持自定义判定规则、时长阈值、调度器:
using System; using System.Reactive.Linq; using System.Reactive.Concurrency; public static class ObservableExtensions { /// <summary> /// 对特殊事件添加持久化判定逻辑,只有持续达到指定时长的特殊事件才会向下游推送 /// </summary> /// <param name="source">源可观测序列</param> /// <param name="isSpecialEvent">特殊事件判定函数,返回true表示当前事件属于特殊类型</param> /// <param name="persistDuration">特殊事件需要持续的最小阈值时长</param> /// <param name="scheduler">时间调度器,默认使用默认调度器,单元测试可传入测试调度器</param> public static IObservable<T> FilterTransientSpecialEvent<T>( this IObservable<T> source, Func<T, bool> isSpecialEvent, TimeSpan persistDuration, IScheduler? scheduler = null) { scheduler ??= Scheduler.Default; return source .Timestamp(scheduler) .Scan( // 状态:首次特殊事件时间、是否已确认特殊事件有效、最后接收的事件 (firstSpecialTime: default(DateTimeOffset?), isSpecialValid: false, lastEvent: default(Timestamped<T>)), (state, currentEvent) => { // 当前为普通事件,重置特殊事件判定状态 if (!isSpecialEvent(currentEvent.Value)) { return (null, false, currentEvent); } // 特殊事件已经过有效性确认,直接保留确认状态 if (state.isSpecialValid) { return (state.firstSpecialTime, true, currentEvent); } // 首次收到特殊事件,记录起始时间 if (state.firstSpecialTime == null) { return (currentEvent.Timestamp, false, currentEvent); } // 判定特殊事件是否达到持续时长阈值 var isValid = currentEvent.Timestamp - state.firstSpecialTime >= persistDuration; return (state.firstSpecialTime, isValid, currentEvent); }) .Where(state => { // 普通事件直接放行 if (!isSpecialEvent(state.lastEvent.Value)) { return true; } // 特殊事件只有确认有效后放行 return state.isSpecialValid; }) .Select(state => state.lastEvent.Value); } }
使用示例
和提供的测试场景对齐,特殊事件判定规则为值等于999,阈值为1.5秒:
// 构造测试数据源 var oneObservable = Observable.Interval(TimeSpan.FromSeconds(1)).Select(_ => 999); var intObservable = Observable.Interval(TimeSpan.FromSeconds(1)).Select(value => (int)value); var myObservable = intObservable.Take(4).Concat(oneObservable.Take(3)).Repeat(); // 调用自定义运算符 var result = myObservable.FilterTransientSpecialEvent( isSpecialEvent: v => v == 999, persistDuration: TimeSpan.FromSeconds(1.5) ); // 订阅输出 result.Subscribe(v => Console.WriteLine($"收到事件:{v}"));
效果说明
- 当特殊事件(如示例中的999)首次出现后1.5秒内没有普通事件插入,达到阈值后所有后续特殊事件都会直接透传
- 若特殊事件首次出现后1.5秒内出现了普通事件,会重置判定状态,后续再次出现特殊事件会重新计算时长
- 所有普通事件全程不受影响,直接透传
内容的提问来源于stack exchange,提问作者Gimly
相关产品推荐
相关产品推荐

