将请求响应方法改造为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风格时,遇到三个核心问题:
核心问题
- 错误处理逻辑误解:认为外层
try/catch块不必要,仅靠.Subscribe即可处理错误; - 超时机制失效:请求超时时未正确返回
Result<IEnumerable<TItemArgs>>.FromError(new TimeoutError(...)); - 错误流合并失败:推送至
_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
相关产品推荐
相关产品推荐

