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

如何在异步方法中等待委托触发特定条件?

解决方案:使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 11:22:24