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

如何实现过滤瞬时特殊事件、仅推送持续指定时长同类事件的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.23 22:06:01