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

如何在调用Publish前将CancellationToken接入IObservable流

如何在调用Publish方法前将CancellationToken接入IObservable冷管道

核心需求与约束

  • 目标:在对现有IObservable管道调用Publish方法(即管道转换为IConnectableObservable)之前,将取消令牌接入管道逻辑
  • 强制约束:取消逻辑必须作为冷可观察对象(cold observable)管道的内置组成部分,在订阅流程完成前就生效。如果无此约束,直接将CancellationToken传入IObservable的Subscribe、RunAsync、ToTask等方法即可实现取消,无需额外改造
  • 疑问:该场景是否存在官方推荐的标准实现模式?

初步实现方案

目前参考通用建议,采用TakeUntil操作符实现取消逻辑:监听取消令牌的触发信号,信号到达时抛出OperationCanceledException终止整个管道。示例代码如下:

using System.Reactive.Linq;
using System.Reactive.Threading.Tasks;

async Task Test(CancellationToken token)
{
    var publishedSequence = Observable
        .Interval(TimeSpan.FromMilliseconds(100))
        .Do(n => Console.WriteLine($"Emitting: {n}"))
        .Skip(3)
        .TakeUntil(
            Observable.Create<long>(
                observer => token.Register(
                    (_, token) => observer.OnError(new OperationCanceledException(token)),
                    null)))
        .Finally(() => Console.WriteLine($"Finally"))
        .Publish();

    using var subscription = publishedSequence.Subscribe(
            onNext: n => Console.WriteLine($"OnNext: {n}"),
            onError: e => Console.WriteLine($"OnError: {e}"),
            onCompleted: () => Console.WriteLine("OnCompleted"));

    using var connection = publishedSequence.Connect();
    await publishedSequence.ToTask();
}

var cts = new CancellationTokenSource(1000);
await Test(cts.Token);

代码运行输出如下:

Emitting: 0
Emitting: 1
Emitting: 2
Emitting: 3
OnNext: 3
Emitting: 4
OnNext: 4
Emitting: 5
OnNext: 5
Emitting: 6
OnNext: 6
Emitting: 7
OnNext: 7
Emitting: 8
OnNext: 8
OnError: System.OperationCanceledException: The operation was canceled.
Finally

已知问题与补充验证

  • 除上述TakeUntil实现外,已开发自定义WithCancellation操作符原型,核心逻辑为透传IObservable事件的同时并行监听取消信号,但更倾向于使用官方标准实现,避免自定义操作符的维护成本
  • 更新:上述直接在token.Register回调中调用observer.OnError的TakeUntil写法存在竞态条件,自研的WithCancellation实现不会触发该问题
  • 更新补充:如果将TakeUntil的参数替换为.TakeUntil(Task.Delay(Timeout.Infinite, token).ToObservable())的写法,同样不会复现前述竞态问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 02:01:24