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

将请求响应方法改造为Rx风格并解决三类实现问题

Rx风格重构Socket请求-响应模型的问题与修复

背景

这是一种基于Socket/WebSocket的请求-响应模型(类似HTTP),通过请求ID匹配对应的响应,工作流程如下:

  • 订阅Error、ContractDetails(_itemObservable)和ContractDetailsEnd(_itemEndObservable)可观察对象;
  • 匹配请求ID与响应ID,通过.OnNext(...)推送消息至ContractDetails;无更多消息时推送ContractDetailsEnd触发.OnCompleted;
  • 清理操作:取消订阅即释放可观察对象。

最初实现使用List<TItemArgs>、CancellationTokenSource和TaskCompletionSource<TItemArgs>,纯Rx风格可完全省去这些冗余代码并精简行数。尝试改造为Rx风格时,遇到三个核心问题:

核心问题

  1. 错误处理逻辑误解:认为外层try/catch块不必要,仅靠.Subscribe即可处理错误;
  2. 超时机制失效:请求超时时未正确返回Result<IEnumerable<TItemArgs>>.FromError(new TimeoutError(...));
  3. 错误流合并失败:推送至_errorSubject的消息无法正确转换为RemoteError返回,.Merge(errorMessages.Any(_ => false))逻辑无效。

选择ReplaySubject而非AsyncSubject的原因:需返回所有响应值,ReplaySubject可记录历史值,确保订阅者收到完整且一致的序列。

初始实现代码

public async ValueTask<Result<IEnumerable<TItemArgs>>> ExecuteAsync(Action<int> action)
{
    var requestId = _client.GetNextRequestId();

    var data = new List<TItemArgs>();
    var cts = new CancellationTokenSource(_timeout);
    var tcs = new TaskCompletionSource<IEnumerable<TItemArgs>>();

    cts.Token.Register(() =>
    {
        tcs.TrySetCanceled();
    }, false);

    void OnError(ErrorData msg)
    {
        tcs.SetException(new IBClientException(msg.RequestId, msg.Code, msg.Message, msg.AdvancedOrderRejectJson));
    }

    void OnDetails(TItemArgs item)
    {
        data.Add(item);
    }

    void OnDetailsEnd(TItemListEndArgs item)
    {
        tcs.TrySetResult(data);
    }

    var disposable = new CompositeDisposable();
    
    _client.Error
        .Where(e => HasRequestId && e.RequestId == requestId)
        .Subscribe(OnError)
        .DisposeWith(disposable);

    _itemObservable
        .Where(item => MatchRequest(item, _itemRequestIdExtractor, requestId))
        .Subscribe(OnDetails)
        .DisposeWith(disposable);

    _itemEndObservable
        .Where(item => MatchRequest(item, _itemListEndRequestIdExtractor, requestId))
        .Subscribe(OnDetailsEnd)
        .DisposeWith(disposable);

    action(requestId);

    try
    {
        await tcs.Task.ContinueWith(x =>
        {
            disposable.Dispose();
            cts.Dispose();
        }, TaskContinuationOptions.RunContinuationsAsynchronously);
        
        return Result<IEnumerable<TItemArgs>>.FromSuccess(tcs.Task.Result);
    }
    catch (Exception e)
    {
        return Result<IEnumerable<TItemArgs>>.FromError(new RemoteError(e.Message, null));
    }
}

Rx风格尝试代码

using System.Reactive.Concurrency;
using System.Reactive.Linq;
using System.Reactive.Subjects;

var client = new IBClient();

var result = await client.GetContractDetailsAsync();

if (result.Success)
{
    foreach (var item in result.Data!)
    {
        Console.WriteLine($"RequestId: {item.RequestId} | Data: {item.ContractDetails}");
    }
}

public sealed class IBClient
{
    private int _nextValidId;
    
    private readonly Subject<ErrorData> _errorSubject = new();
    public IObservable<ErrorData> Error => _errorSubject.AsObservable();
    
    private readonly Subject<ContractDetailsData> _contractDetailsSubject = new();
    public IObservable<ContractDetailsData> ContractDetails => _contractDetailsSubject.AsObservable();

    private readonly Subject<RequestEndData> _contractDetailsEndSubject = new();
    public IObservable<RequestEndData> ContractDetailsEnd => _contractDetailsEndSubject.AsObservable();
    
    public ValueTask<Result<IEnumerable<ContractDetailsData>>> GetContractDetailsAsync()
    {
        return new PendingRequest<ContractDetailsData, RequestEndData>(
                this,
                ContractDetails,
                ContractDetailsEnd,
                e => e.RequestId,
                e => e.RequestId)
            .ExecuteAsync(reqId => TestCall(reqId));
    }

    private void TestCall(int requestId)
    {
        _contractDetailsSubject.OnNext(new ContractDetailsData(requestId, "hey from test call"));
        _contractDetailsSubject.OnNext(new ContractDetailsData(requestId, "hey two"));

        // TODO: Errors doesn't seem to work
        // _errorSubject.OnNext(new ErrorData(requestId, 123, "Error happened", ""));

        _contractDetailsEndSubject.OnNext(new RequestEndData(requestId));

        // There shouldn't be matched.
        _contractDetailsSubject.OnNext(new ContractDetailsData(123, "fake ones, so we know it works"));
        _contractDetailsEndSubject.OnNext(new RequestEndData(123));
    }
    
    public int GetNextRequestId()
    {
        return Interlocked.Increment(ref _nextValidId);
    }
}

public sealed class PendingRequest<TItemArgs, TItemListEndArgs>
{
    private readonly TimeSpan _timeout = TimeSpan.FromSeconds(2);
    private readonly IBClient _client;

    private readonly IObservable<TItemArgs> _itemObservable;
    private readonly IObservable<TItemListEndArgs> _itemEndObservable;
    private readonly Func<TItemArgs, int>? _itemRequestIdExtractor;
    private readonly Func<TItemListEndArgs, int>? _itemListEndRequestIdExtractor;

    public PendingRequest(
        IBClient client,
        IObservable<TItemArgs> itemObservable,
        IObservable<TItemListEndArgs> itemEndObservable,
        Func<TItemArgs, int>? itemRequestIdExtractor = null,
        Func<TItemListEndArgs, int>? itemListEndRequestIdExtractor = null)
    {
        _client = client;

        _itemObservable = itemObservable;
        _itemEndObservable = itemEndObservable;

        _itemRequestIdExtractor = itemRequestIdExtractor;
        _itemListEndRequestIdExtractor = itemListEndRequestIdExtractor;
    }

    private bool HasRequestId => _itemRequestIdExtractor != null && _itemListEndRequestIdExtractor != null;

    public async ValueTask<Result<IEnumerable<TItemArgs>>> ExecuteAsync(Action<int> action, IScheduler? scheduler = null)
    {
        scheduler ??= ImmediateScheduler.Instance;

        var requestId = _client.GetNextRequestId();
        var results = new ReplaySubject<TItemArgs>();

        try
        {
            var errorMessages = _client.Error
                .Where(e => HasRequestId && e.RequestId == requestId);

            using (_itemObservable
                       .Where(item => MatchRequest(item, _itemRequestIdExtractor, requestId))
                       .ObserveOn(scheduler)
                       .Subscribe(results))
            using (_itemEndObservable
                       .Any(item => MatchRequest(item, _itemListEndRequestIdExtractor, requestId))
                       .Merge(errorMessages.Any(_ => false)) // TODO: ???
                       .ObserveOn(scheduler)
                       .Subscribe(_ => results.OnCompleted()))
            {
                action(requestId);

                // Don't want an Exception thrown if there result list is empty
                await results.DefaultIfEmpty();

                return Result<IEnumerable<TItemArgs>>.FromSuccess(results.ToEnumerable());
            }
        }
        catch (Exception ex)
        {
            return Result<IEnumerable<TItemArgs>>.FromError(new RemoteError(ex.Message, null));
        }
    }
    
    private bool MatchRequest<T>(T item, Func<T, int>? idExtractor, int id)
    {
        return !HasRequestId || (idExtractor != null && idExtractor(item) == id);
    }
}

public sealed class ContractDetailsData
{
    public ContractDetailsData(int requestId, string contractDetails)
    {
        RequestId = requestId;
        ContractDetails = contractDetails;
    }

    public int RequestId { get; }

    public string ContractDetails { get; }
}

public sealed class ErrorData
{
    public ErrorData(int requestId, int code, string message, string advancedOrderRejectJson)
    {
        RequestId = requestId;
        Code = code;
        Message = message;
        AdvancedOrderRejectJson = advancedOrderRejectJson;
    }

    public int RequestId { get; }
    public int Code { get; }
    public string Message { get; }
    public string AdvancedOrderRejectJson { get; }
}

public sealed class RequestEndData
{
    public RequestEndData(int requestId)
    {
        RequestId = requestId;
    }

    public int RequestId { get; }
}

public class IBClientException : Exception
{
    public IBClientException(int requestId, int errorCode, string message, string advancedOrderRejectJson)
        : base(message)
    {
        RequestId = requestId;
        ErrorCode = errorCode;
        AdvancedOrderRejectJson = advancedOrderRejectJson;
    }

    public IBClientException(string err)
        : base(err)
    {
    }

    public IBClientException(Exception e)
    {
        Exception = e;
    }

    public int RequestId { get; }
    public int ErrorCode { get; }
    public string? AdvancedOrderRejectJson { get; }
    public Exception? Exception { get; }
}

public abstract record Error(int? Code, string Message, object? Data);

public record RemoteError : Error
{
    public RemoteError(string message, object? data) : base(null, message, data)
    {
    }

    public RemoteError(int? code, string message, object? data) : base(code, message, data)
    {
    }
}

public record Result<T>(bool Success, T? Data, Error? Error)
{
    public Result(T data) : this(true, data, default)
    {
    }

    public Result(Error error) : this(false, default, error)
    {
    }

    public static Result<T> FromSuccess(T data)
    {
        return new Result<T>(data);
    }

    public static Result<T> FromError<TError>(TError error) where TError : Error
    {
        return new Result<T>(error);
    }
}

public static class DisposableExtensions
{
    public static T DisposeWith<T>(this T disposable, ICollection<IDisposable> collection)
        where T : IDisposable
    {
        ArgumentNullException.ThrowIfNull(disposable);
        ArgumentNullException.ThrowIfNull(collection);

        collection.Add(disposable);
        return disposable;
    }
}

问题修复方案

1. 错误处理优化:用Rx操作符替代外层try/catch

外层try/catch无法覆盖所有Rx流内的错误,需用Rx的Catch操作符统一处理错误,并将异常转换为Result的错误类型。同时,错误流需要触发序列终止,避免订阅者等待。

2. 超时机制修复

在合并后的序列上添加Timeout操作符,超时后抛出TimeoutException,再通过Catch转换为TimeoutError(需先定义该错误类型)。

3. 错误流合并修复

将错误流映射为异常,然后合并到主序列中,或者将错误信号作为终止信号触发错误处理。替换原来的errorMessages.Any(_ => false)为errorMessages.Select(_ => throw new IBClientException(...)),或者将错误流转换为终止信号并触发错误。

修复后的核心ExecuteAsync方法

public async ValueTask<Result<IEnumerable<TItemArgs>>> ExecuteAsync(Action<int> action, IScheduler? scheduler = null)
{
    scheduler ??= ImmediateScheduler.Instance;
    var requestId = _client.GetNextRequestId();

    // 过滤当前请求的数据流
    var filteredItems = _itemObservable
        .Where(item => MatchRequest(item, _itemRequestIdExtractor, requestId))
        .ObserveOn(scheduler);

    // 过滤当前请求的结束信号
    var endSignal = _itemEndObservable
        .Where(item => MatchRequest(item, _itemListEndRequestIdExtractor, requestId))
        .Take(1)
        .ObserveOn(scheduler);

    // 过滤当前请求的错误流,并转换为异常
    var errorSignal = _client.Error
        .Where(e => HasRequestId && e.RequestId == requestId)
        .Take(1)
        .Select(e => throw new IBClientException(e.RequestId, e.Code, e.Message, e.AdvancedOrderRejectJson))
        .ObserveOn(scheduler);

    // 合并数据流、结束信号、错误信号,设置超时
    var sequence = filteredItems
        .TakeUntil(endSignal.Merge(errorSignal))
        .Timeout(_timeout, scheduler)
        .Catch((Exception ex) => Observable.Throw<TItemArgs>(ex));

    try
    {
        action(requestId);
        var items = await sequence.ToList();
        return Result<IEnumerable<TItemArgs>>.FromSuccess(items);
    }
    catch (TimeoutException)
    {
        return Result<IEnumerable<TItemArgs>>.FromError(new RemoteError("请求超时", null));
    }
    catch (IBClientException ex)
    {
        return Result<IEnumerable<TItemArgs>>.FromError(new RemoteError(ex.ErrorCode, ex.Message, ex.AdvancedOrderRejectJson));
    }
    catch (Exception ex)
    {
        return Result<IEnumerable<TItemArgs>>.FromError(new RemoteError(ex.Message, null));
    }
}

修复说明

  • 去掉了ReplaySubject,直接用ToList()收集所有结果,更简洁;
  • 用TakeUntil监听结束信号和错误信号,自动终止数据流;
  • Timeout操作符直接应用在合并后的序列上,确保超时触发;
  • 错误流直接转换为异常,通过Catch和外层try/catch处理并转换为Result类型;
  • 无需手动管理订阅,Rx会自动处理资源释放。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 03:01:39