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

复杂Async/Await场景下如何实现超时、取消与重试机制?

解决异步分离架构中的超时与取消重试问题

你的核心问题在于当前的TaskCompletionSource(TCS)没有绑定取消机制,重试时只是新增请求,没有清理之前未完成的请求上下文。下面是一套可行的改造方案,从核心的Bot类到重试逻辑逐步调整:

1. 给Bot类添加取消支持

首先要让Bot的每个请求都能被取消,我们需要把TaskCompletionSource和CancellationTokenSource(CTS)绑定起来,同时记录当前请求的唯一标识,方便后续从队列移除:

class Bot {
    private TaskCompletionSource<object> _taskCompSource;
    private CancellationTokenSource _cts;
    // 记录当前请求的唯一标识,用于在Service队列中定位
    public Guid CurrentRequestId { get; private set; }
    public string Id { get; set; } // 新增Bot标识,用于关联Service队列

    private void StartTask() {
        // 先取消之前未完成的请求(如果有)
        _cts?.Cancel();
        _cts?.Dispose();
        _cts = new CancellationTokenSource();
        _taskCompSource = new TaskCompletionSource<object>(_cts.Token);
        CurrentRequestId = Guid.NewGuid();
    }

    private async Task<T> ResultPromise<T>() {
        try {
            var result = await _taskCompSource.Task.ConfigureAwait(false);
            return result != null ? (T)result : default(T);
        } catch (OperationCanceledException) {
            // 取消时返回默认值或抛出异常,根据业务需求调整
            return default(T);
        }
    }

    public Task<string> GetSomething() {
        StartTask();
        // 把请求ID和BotID传入Request,方便后续定位
        QueueRequest?.Invoke(this, new Request(Id, CurrentRequestId, ...));
        return ResultPromise<string>();
    }

    public void CancelCurrentRequest() {
        _cts?.Cancel();
        _taskCompSource?.TrySetCanceled();
    }

    public void Process(Reply reply) {
        // 检查当前请求是否已取消,避免处理过期回复
        if (_cts?.IsCancellationRequested ?? false) return;
        _taskCompSource.SetResult("success");
    }
}

2. 让Service支持移除指定请求

修改Service的队列逻辑,允许根据BotID和请求ID移除未发送的请求:

class Service {
    private Dictionary<string, Queue<Request>> _queues;

    // 新增方法:移除指定Bot队列中的特定请求
    public void CancelRequestForBot(string botId, Guid requestId) {
        if (_queues.TryGetValue(botId, out var queue)) {
            // 过滤掉指定ID的请求,替换原队列
            var filteredQueue = new Queue<Request>(queue.Where(r => r.RequestId != requestId));
            _queues[botId] = filteredQueue;
        }
    }

    public void QueueRequestForSending(Request request) {
        var queue = GetQueueForBot(request.BotId);
        queue.Enqueue(request);
    }

    private async void HttpCallback(IAsyncResult result) {
        while (true) {
            // 处理回复逻辑...
            // 发送请求前检查有效性
            foreach (var botId in _queues.Keys.ToList()) {
                if (_queues[botId].Count == 0) continue;
                
                var pendingRequest = _queues[botId].Peek();
                var bot = GetBot(pendingRequest.BotId);
                // 如果请求已被取消,直接出队丢弃
                if (bot._cts.IsCancellationRequested) {
                    _queues[botId].Dequeue();
                    continue;
                }
                // 发送请求逻辑...
            }
            // 保持连接的空格发送逻辑...
        }
    }
}

3. 重写Retry扩展方法,加入取消逻辑

现在的重试方法需要在每次超时后,主动取消上一次的请求,再发起新请求:

static class Extensions {
    public static async Task<T> WithRetry<T>(this Bot bot, Func<Task<T>> taskFactory, int retries = 3, int timeout = 10000) {
        for (int attempt = 0; attempt <= retries; attempt++) {
            // 每次重试前,取消上一次未完成的请求
            bot.CancelCurrentRequest();
            // 通知Service移除队列中的旧请求
            BotManager.Instance.GetService().CancelRequestForBot(bot.Id, bot.CurrentRequestId);

            var task = taskFactory();
            var timeoutTask = Task.Delay(timeout);

            var completedTask = await Task.WhenAny(task, timeoutTask).ConfigureAwait(false);

            if (completedTask == task) {
                // 请求成功完成,返回结果
                return await task.ConfigureAwait(false);
            }

            // 最后一次重试仍超时,抛出异常或返回默认值
            if (attempt == retries) {
                throw new TimeoutException($"请求超时,已完成{retries+1}次尝试");
            }
        }
        return default(T);
    }
}

4. 调用方式调整

改用带Bot实例的扩展方法调用,确保重试时能正确取消旧请求:

class CLI { 
    async void Main() { 
        var bot = BotManager.Instance.GetBot();
        try {
            var response = await bot.WithRetry(() => bot.GetSomething(), retries: 3, timeout: 10000);
            // 处理成功响应
        } catch (TimeoutException ex) {
            // 处理超时逻辑
        }
    } 
}

关键要点总结

  • 绑定TCS与CTS:每个请求对应独立的TaskCompletionSource和CancellationTokenSource,确保超时后能终止等待。
  • 请求标识追踪:给每个Request添加唯一ID和BotID,方便在Service队列中定位并移除过期请求。
  • 重试时主动清理:每次重试前先取消Bot的当前请求,再通知Service清理队列,避免无效请求堆积。
  • 资源释放:取消后及时释放CancellationTokenSource资源,防止内存泄漏。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 05:16:02