Reactive Extensions C#:长订阅逻辑执行时如何跳过消息?
Absolutely, Rx.NET has built-in constructs to handle this exact scenario—where you want to skip incoming messages while a long-running processing task is active, then pick up the latest (or next) message once processing completes.
The core idea is to signal when processing finishes, then use that signal to "gate" which messages get processed. Here are two clean, built-in approaches:
Approach 1: Using Sample with a Completion Signal
Sample emits the latest value from the source observable whenever a trigger observable emits a value. We can use this to only process messages once our processing task has finished.
using System.Reactive.Subjects; using System.Reactive.Linq; // Create a subject to signal when processing is complete var processingComplete = new BehaviorSubject<Unit>(Unit.Default); // Seed initial value to start processing IObservable<int> stream = ...; stream .Sample(processingComplete) // Only take the latest message when processing completes .SelectMany(msg => // Wrap processing in an observable to track completion Observable.FromAsync(() => ProcessMessage(msg)) .Do(_ => processingComplete.OnNext(Unit.Default)) // Signal completion after processing ) .SubscribeOn(TaskPoolScheduler.Default) .Subscribe( _ => {}, // No need to handle result unless required ex => Console.WriteLine($"Error: {ex.Message}"), () => processingComplete.OnCompleted() );
How this works:
- The
BehaviorSubjectstarts with an initialUnitvalue, so the first message is processed immediately. - While
ProcessMessageruns, any incoming messages are ignored bySampleuntil processing finishes. - When processing completes, we emit a signal to
processingComplete, which triggersSampleto take the latest message from the stream (skipping all intermediate ones) and start processing it.
Approach 2: Using Buffer with a Closing Selector
Buffer collects messages until a closing selector observable emits, then emits the collected buffer. We can take the last message from each buffer to skip intermediates.
using System.Reactive.Subjects; using System.Reactive.Linq; var processingComplete = new BehaviorSubject<Unit>(Unit.Default); IObservable<int> stream = ...; stream .Buffer(() => processingComplete) // Collect messages until processing completes .Where(buffer => buffer.Any()) // Skip empty buffers (if no messages arrived during processing) .Select(buffer => buffer.Last()) // Take only the latest message from the buffer .SelectMany(msg => Observable.FromAsync(() => ProcessMessage(msg)) .Do(_ => processingComplete.OnNext(Unit.Default)) ) .SubscribeOn(TaskPoolScheduler.Default) .Subscribe( _ => {}, ex => Console.WriteLine($"Error: {ex.Message}"), () => processingComplete.OnCompleted() );
Key Notes:
- Both approaches ensure only one instance of
ProcessMessageruns at a time. - If you prefer to take the first message that arrives after processing (instead of the latest), replace
buffer.Last()withbuffer.First()—though this is less common for backpressure scenarios like yours. - Always clean up subjects by calling
OnCompleted()when done to avoid memory leaks.
Adjustment for Synchronous Processing
If ProcessMessage is synchronous, simplify the code by removing Observable.FromAsync and handling completion directly:
stream .Sample(processingComplete) .Subscribe( msg => { ProcessMessage(msg); processingComplete.OnNext(Unit.Default); }, ex => Console.WriteLine($"Error: {ex.Message}"), () => processingComplete.OnCompleted() );
内容的提问来源于stack exchange,提问作者stackh34p

