复杂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
相关产品推荐
相关产品推荐

