如何在异步方法中等待委托触发特定条件?
解决方案:使用
TaskCompletionSource实现异步等待 在异步场景下,阻塞式的同步等待(比如ManualResetEvent.WaitOne())会占用线程,导致WebSocket的消息回调无法正常执行。正确的做法是使用**TaskCompletionSource<T>**,它能在异步环境下安全地等待特定条件满足,且不会阻塞线程。
修改核心代码
1. 调整BuyAndSetLimitSellOrder方法,引入TaskCompletionSource
private async Task BuyAndSetLimitSellOrder(string pair, CancellationToken cancellationToken) { var buyOrder = await BuyAndWriteOrderToDb(pair); var sellOrder = await SetSellAndWriteOrderToDb(pair, buyOrder); // 创建异步等待源,指定异步续期避免死锁 var completionSource = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously); (Func<Task> pingDelegate, Func<Task> closeDelegate)? subscription = null; // 注册取消令牌回调,防止永久等待 using var cancellationRegistration = cancellationToken.Register(async () => { completionSource.TrySetCanceled(); if (subscription?.closeDelegate != null) await subscription.Value.closeDelegate(); }); // 订阅用户数据流,传入completionSource用于触发等待完成 subscription = await _binanceService.SubscribeToUserData(async userData => { await HandleReceivedMessage(userData, sellOrder.Id, completionSource); }, cancellationToken); try { // 异步等待直到满足条件或被取消 await completionSource.Task; } finally { // 清理资源:关闭WebSocket并销毁监听密钥 if (subscription?.closeDelegate != null) await subscription.Value.closeDelegate(); } }
2. 修改HandleReceivedMessage方法,满足条件时触发等待完成
private async Task HandleReceivedMessage(string userData, long targetSellOrderId, TaskCompletionSource<bool> completionSource) { // 解析Binance用户数据消息(根据实际结构调整) var userDataEvent = System.Text.Json.JsonSerializer.Deserialize<BinanceOrderUpdateEvent>(userData); // 判断是否满足目标条件:比如订单成交、包含特定字符串等 if (userDataEvent?.EventType == "orderTradeUpdate" && userDataEvent.OrderId == targetSellOrderId && userDataEvent.OrderStatus == "FILLED") { // 满足条件,标记等待完成 completionSource.TrySetResult(true); } // 其他消息处理逻辑(如日志、状态更新等) await Task.CompletedTask; }
3. 优化SubscribeToUserData方法,返回清理资源的委托
原方法只返回了ping委托,无法关闭连接,修改后返回ping和关闭的组合委托,方便后续清理:
public async Task<(Func<Task> Ping, Func<Task> Close)> SubscribeToUserData(Func<string, Task> messageHandler, CancellationToken cancellationToken) { string response = await _userDataStreams.CreateSpotListenKey(); string listenKey = (System.Text.Json.JsonSerializer.Deserialize<CreateSpotListenKeyResponse>(response) ?? throw new FormatException("Invalid AccountInformation string")).ListenKey; UserDataWebSocket websocket = new(listenKey, _configuration.BinanceWssBaseUrl); websocket.OnMessageReceived(messageHandler, CancellationToken.None); await websocket.ConnectAsync(CancellationToken.None); // 定义关闭逻辑:断开WebSocket并销毁监听密钥 async Task CloseStream() { await websocket.DisconnectAsync(CancellationToken.None); await _userDataStreams.DeleteSpotListenKey(listenKey); } return (() => PingSpotListenKey(listenKey), CloseStream); }
关键说明
TaskCompletionSource的优势:它是异步友好的信号机制,await completionSource.Task不会阻塞线程,线程可以继续处理WebSocket的消息回调,避免了阻塞导致的消息无法处理问题。- 资源清理:通过返回的
Close委托,在等待完成或取消时关闭WebSocket并销毁Binance的监听密钥,避免资源泄漏。 - 取消令牌处理:注册取消令牌的回调,确保外部取消时能终止等待并清理资源,防止任务永久挂起。
内容的提问来源于stack exchange,提问作者exzzy
相关产品推荐
相关产品推荐

