如何让Reactive Extensions的onNext、onError、onCompleted支持async/await
问题描述
用户实现了IObservable<T>的扩展方法Observe,试图让订阅支持等待,但核心问题是onNext、onError、onCompleted委托不支持await——仅对整个订阅做了包装,且已知有针对onNext的异步方案,但未覆盖onError和onCompleted,需要实现全异步的订阅支持。
现有代码:
public static Task Observe<T>(this IObservable<T> observable, Action<T> onNext, Func<Exception, EventWaitHandle, ExceptionDispatchInfo> onError, Action<EventWaitHandle, CancellationToken> onCompleted, CancellationToken cancellationToken ) where T : class { var compositeDisposables = new CompositeDisposable(); var waitHandle = new ManualResetEvent(false); compositeDisposables.Add(waitHandle); var disposable = observable.Subscribe( // This should be await onNext onNext, // This should be await onError e => onError(e, waitHandle), // This should be await onCompleted () => onCompleted(waitHandle, cancellationToken)); compositeDisposables.Add(disposable); waitHandle.WaitOne(); compositeDisposables.Dispose(); return Task.CompletedTask; }
解决方案
要实现全异步的订阅回调,需从委托类型、异步信号处理、生命周期管理三个核心方向修改:
核心实现代码
public static async Task ObserveAsync<T>(this IObservable<T> observable, Func<T, Task> onNextAsync, Func<Exception, Task> onErrorAsync, Func<CancellationToken, Task> onCompletedAsync, CancellationToken cancellationToken = default) where T : class { using var cts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); var tcs = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously); using var compositeDisposable = new CompositeDisposable(); // 绑定取消令牌,取消时终止订阅并完成任务 compositeDisposable.Add(cts.Token.Register(() => { tcs.TrySetCanceled(cts.Token); })); var disposable = observable.Subscribe( async value => { try { await onNextAsync(value).ConfigureAwait(false); } catch (Exception ex) { cts.Cancel(); tcs.TrySetException(ex); } }, async ex => { try { await onErrorAsync(ex).ConfigureAwait(false); tcs.TrySetResult(true); } catch (Exception callbackEx) { tcs.TrySetException(callbackEx); } finally { cts.Cancel(); } }, async () => { try { await onCompletedAsync(cts.Token).ConfigureAwait(false); tcs.TrySetResult(true); } catch (Exception ex) { tcs.TrySetException(ex); } finally { cts.Cancel(); } }); compositeDisposable.Add(disposable); try { await tcs.Task.ConfigureAwait(false); } finally { compositeDisposable.Dispose(); } }
关键改动说明
- 异步委托签名:将所有回调参数改为
Func<..., Task>类型,允许内部使用await。 - 非阻塞等待:用
TaskCompletionSource替代ManualResetEvent的阻塞等待,适配异步编程模型,避免线程资源浪费。 - 异常与取消处理:捕获每个异步回调的异常并传递到外部任务,绑定取消令牌实现主动终止订阅的逻辑。
- 性能优化:使用
ConfigureAwait(false)避免不必要的同步上下文切换,提升异步执行效率。
建议补充的扩展方法重载
为适配不同使用场景,需补充以下重载:
1. 仅提供onNextAsync的简化重载
public static Task ObserveAsync<T>(this IObservable<T> observable, Func<T, Task> onNextAsync, CancellationToken cancellationToken = default) where T : class { return observable.ObserveAsync( onNextAsync, ex => Task.FromException(ex), _ => Task.CompletedTask, cancellationToken); }
2. 包含onNextAsync和onErrorAsync的重载
public static Task ObserveAsync<T>(this IObservable<T> observable, Func<T, Task> onNextAsync, Func<Exception, Task> onErrorAsync, CancellationToken cancellationToken = default) where T : class { return observable.ObserveAsync( onNextAsync, onErrorAsync, _ => Task.CompletedTask, cancellationToken); }
3. 兼容同步回调的重载
如果需要兼容原有同步回调场景,可自动将同步委托包装为异步:
public static Task ObserveAsync<T>(this IObservable<T> observable, Action<T> onNext, Action<Exception> onError, Action onCompleted, CancellationToken cancellationToken = default) where T : class { return observable.ObserveAsync( value => { onNext(value); return Task.CompletedTask; }, ex => { onError(ex); return Task.CompletedTask; }, _ => { onCompleted(); return Task.CompletedTask; }, cancellationToken); }
内容的提问来源于stack exchange,提问作者xplat
相关产品推荐
相关产品推荐

