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

Reactive Extensions C#:长订阅逻辑执行时如何跳过消息?

Handling Backpressure by Skipping Messages During Processing in Rx.NET

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:

  1. The BehaviorSubject starts with an initial Unit value, so the first message is processed immediately.
  2. While ProcessMessage runs, any incoming messages are ignored by Sample until processing finishes.
  3. When processing completes, we emit a signal to processingComplete, which triggers Sample to 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 ProcessMessage runs at a time.
  • If you prefer to take the first message that arrives after processing (instead of the latest), replace buffer.Last() with buffer.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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 21:32:59