如何为CancellationTokenSource实现滑动过期(重置超时时间)
问题
在C#中使用TaskCompletionSource<T>配合CancellationTokenSource实现超时逻辑时,需要为CancellationTokenSource添加滑动过期功能——也就是能重置超时时间。具体场景是应用发送命令后等待回复,回复时长不固定,要求如果1秒内没有新数据就触发超时,固定超时设置太长,滑动过期更适配。
尝试过自行实现定时器,但发现CancellationTokenSource内部用的是密封定时器,想找合适的实现方式。现有代码示例如下:
TaskCompletionSource<T> source = new TaskCompletionSource<T>(); CancellationTokenSource cancellationTokenSource = new CancellationTokenSource(timeoutMs); cancellationTokenSource.Token.Register(() => source.TrySetCanceled(), useSynchronizationContext: false); // do work and eventually call source.TrySetResult() // what is the right way to prolong CancellationTokenSource timeout? return source.Task;
实现滑动超时的方案
原生CancellationTokenSource没有内置滑动超时的能力,我们可以基于System.Threading.Timer手动实现,既能控制超时重置,又能和TaskCompletionSource联动。
核心思路
- 用可重置的定时器跟踪超时时间
- 收到新数据时调用定时器重置方法,重新开始计时
- 定时器触发时,取消
CancellationTokenSource并完成TaskCompletionSource - 任务完成(成功/失败/超时)时,及时销毁定时器和取消源,避免资源泄漏
适配你场景的改造代码
TaskCompletionSource<T> source = new TaskCompletionSource<T>(); CancellationTokenSource cancellationTokenSource = new CancellationTokenSource(); // 初始化定时器,初始超时1秒 Timer slidingTimer = new Timer(_ => { cancellationTokenSource.TryCancel(); }, null, 1000, Timeout.Infinite); // 取消时触发TaskCompletionSource的取消逻辑 cancellationTokenSource.Token.Register(() => { slidingTimer.Dispose(); source.TrySetCanceled(); }, useSynchronizationContext: false); // 重置超时的方法,收到新数据时调用它 void ResetSlidingTimeout() { // 重新设置定时器,从当前时间开始再等1秒 slidingTimer.Change(1000, Timeout.Infinite); } // do work and eventually call source.TrySetResult() // 收到新数据时执行:ResetSlidingTimeout(); // 任务结束后清理资源 source.Task.ContinueWith(_ => { slidingTimer.Dispose(); cancellationTokenSource.Dispose(); }, TaskContinuationOptions.ExecuteSynchronously); return source.Task;
封装成可复用的工具方法
如果需要在多个场景使用,可以封装成通用方法:
public class SlidingTimeoutTask<T> { public Task<T> Task { get; } private readonly Timer _timer; private readonly CancellationTokenSource _cts; private readonly TaskCompletionSource<T> _tcs; public SlidingTimeoutTask(int slidingTimeoutMs) { _tcs = new TaskCompletionSource<T>(); _cts = new CancellationTokenSource(); _timer = new Timer(_ => { _cts.TryCancel(); }, null, slidingTimeoutMs, Timeout.Infinite); _cts.Token.Register(() => { _timer.Dispose(); _tcs.TrySetCanceled(); }, useSynchronizationContext: false); // 任务结束自动清理 _tcs.Task.ContinueWith(_ => { _timer.Dispose(); _cts.Dispose(); }, TaskContinuationOptions.ExecuteSynchronously); Task = _tcs.Task; } // 重置超时 public void ResetTimeout() { _timer.Change(1000, Timeout.Infinite); } // 设置任务结果 public bool TrySetResult(T result) { return _tcs.TrySetResult(result); } // 设置任务异常 public bool TrySetException(Exception ex) { return _tcs.TrySetException(ex); } } // 使用示例 var slidingTask = new SlidingTimeoutTask<YourResultType>(1000); // 收到新数据时调用 slidingTask.ResetTimeout(); // 拿到结果时调用 slidingTask.TrySetResult(result); return slidingTask.Task;
内容的提问来源于stack exchange,提问作者Alexander Selishchev
相关产品推荐
相关产品推荐

