Rx Observable技术需求:实现按固定间隔获取源最新值直到源完成
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 insideConcat, 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,
Switchcancels the previous repeat sequence and starts fresh with the new value. - Once the source completes after emitting
F, the operator stops immediately—soFnever gets emitted, just like your example shows.
Key Notes
- Both implementations handle error signals correctly: if the source observable throws an error, the custom
TakeLatestoperator will propagate that error immediately. - Use
Publishin the first implementation if your source has side effects, to avoid redundant work from multiple subscriptions.
内容的提问来源于stack exchange,提问作者Bogey

