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

Rx.NET中如何用触发器获取源Observable最新值?避免Buffer内存低效问题

你可以通过组合Rx.NET的现有运算符来实现这个需求,无需自定义新运算符,还能避免Buffer带来的内存浪费。核心思路是追踪源序列的最新值,同时记录自上次触发器触发后源是否有过更新,具体实现如下:

实现代码

// 发布源序列,避免多次订阅导致重复执行
var publishedSource = dataObservable.Publish();

// 合并源序列与触发器序列,统一处理状态
var merged = publishedSource
    .Select(val => new { IsTrigger = false, Value = val })
    .Merge(triggerObservable.Select(_ => new { IsTrigger = true, Value = default(Data) }))
    .Publish();

// 用Scan维护状态:保存源的最新值、标记是否有未发送的更新
var result = merged
    .Scan(
        new { LatestValue = default(Data), HasUnsentUpdate = false },
        (state, currentItem) =>
        {
            if (!currentItem.IsTrigger)
            {
                // 源产生新值:更新最新值,标记存在未发送更新
                return new { LatestValue = currentItem.Value, HasUnsentUpdate = true };
            }
            else
            {
                // 触发器触发:保留最新值,重置未发送标记
                return new { state.LatestValue, HasUnsentUpdate = false };
            }
        }
    )
    // 关联状态与原事件,筛选符合条件的触发器触发行为
    .Zip(merged, (state, item) => (State: state, Item: item))
    .Where(pair => pair.Item.IsTrigger && pair.State.HasUnsentUpdate)
    .Select(pair => pair.State.LatestValue);

// 订阅结果序列
result.Subscribe(_ => Trigger());

// 连接发布的序列,启动数据流
publishedSource.Connect();
merged.Connect();

逻辑说明

  1. Publish共享序列:用Publish包装源和合并后的序列,确保所有订阅共享同一数据流,避免重复计算和资源浪费。
  2. 状态追踪:通过Scan运算符维护两个核心状态:
    • LatestValue:始终存储源序列的最新值
    • HasUnsentUpdate:标记自上次触发器触发后,源是否产生过新值
  3. 筛选有效触发:通过Zip将状态与原合并序列关联,仅在触发器触发且存在未发送更新时,输出源的最新值。

这种方式仅保留源的最新值,不会像Buffer那样存储所有中间值,内存效率大幅提升。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 06:36:49