如何在调用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
相关产品推荐
相关产品推荐

