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();
逻辑说明
Publish共享序列:用Publish包装源和合并后的序列,确保所有订阅共享同一数据流,避免重复计算和资源浪费。- 状态追踪:通过
Scan运算符维护两个核心状态:LatestValue:始终存储源序列的最新值HasUnsentUpdate:标记自上次触发器触发后,源是否产生过新值
- 筛选有效触发:通过
Zip将状态与原合并序列关联,仅在触发器触发且存在未发送更新时,输出源的最新值。
这种方式仅保留源的最新值,不会像Buffer那样存储所有中间值,内存效率大幅提升。
内容的提问来源于stack exchange,提问作者Tomer B
相关产品推荐
相关产品推荐

