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

如何让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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 12:32:49