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

Rx Observable技术需求:实现按固定间隔获取源最新值直到源完成

Custom TakeLatest Operator for Rx.NET

Great question! There isn't a built-in Rx operator that exactly matches this behavior out of the box, but you can easily assemble it using a combination of standard Rx operators. Let's break this down to meet all three of your requirements:

Step-by-Step Implementation

Here's a clean implementation that aligns with your written requirements:

public static IObservable<T> TakeLatest<T>(this IObservable<T> input, TimeSpan interval)
{
    // Use Publish to avoid multiple subscriptions to the source observable
    return input.Publish(publishedSource =>
        // 1. Immediately emit the first element
        publishedSource.Take(1)
            .Concat(
                // 2. After the first element, emit the latest value at fixed intervals
                Observable.Interval(interval)
                    .WithLatestFrom(publishedSource, (_, latestValue) => latestValue)
                    // 3. Stop as soon as the source completes or errors
                    .TakeUntil(publishedSource.IgnoreElements().Materialize())
            )
    );
}

How It Works

Let's walk through each part to see how it satisfies your requirements:

  • Publish(publishedSource): Ensures we only subscribe to the source observable once, preventing duplicate emissions or side effects from multiple subscriptions.
  • publishedSource.Take(1): Immediately emits the first element from the source, fulfilling your first requirement.
  • Observable.Interval(interval): Generates a signal at your specified fixed interval. Since this is inside Concat, it only starts emitting after the first element has been emitted, aligning with your second requirement of starting the interval clock after the first output.
  • WithLatestFrom(...): Each time the interval signal fires, it grabs the most recent value from the source observable and emits it.
  • TakeUntil(...): Listens for the source observable's completion or error signal. As soon as the source finishes, this stops the interval-based emissions, ensuring the operator follows the source's lifecycle (your third requirement).

Matching Your Marble Diagram

Your sample output suggests you might want to emit new elements immediately and repeat the latest element at intervals (instead of only emitting at fixed interval points). Here's an alternative implementation that perfectly matches your marble diagram:

public static IObservable<T> TakeLatest<T>(this IObservable<T> input, TimeSpan interval)
{
    return input.Select(value =>
            // Emit the new value immediately, then repeat it at intervals
            Observable.Return(value)
                .Concat(Observable.Interval(interval).Select(_ => value))
        )
        // Switch to the latest sequence whenever a new value arrives
        .Switch()
        // Stop when the source completes
        .TakeUntil(input.IgnoreElements().Materialize());
}

For your example with a 2-time-unit interval:

  • Source sequence: A(1) → B(2) → C(3) → D(4) → E(5) → F(6) → Complete
  • Output sequence: A(1) → B(2) → D(4) → D(6) → E(5) → E(7) → Complete

This works because:

  • Every new value triggers an immediate emission, then starts repeating that value at your interval.
  • When a new value arrives, Switch cancels the previous repeat sequence and starts fresh with the new value.
  • Once the source completes after emitting F, the operator stops immediately—so F never gets emitted, just like your example shows.

Key Notes

  • Both implementations handle error signals correctly: if the source observable throws an error, the custom TakeLatest operator will propagate that error immediately.
  • Use Publish in the first implementation if your source has side effects, to avoid redundant work from multiple subscriptions.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 11:02:34