多实例微服务:如何实现跨实例任务取消?
解决方案思路
核心问题拆解
你面临的核心矛盾是:分布式环境下要跨实例控制任务取消,但CancellationToken是内存内对象,没法直接序列化存入分布式缓存。而你需要一套能兼容本地内存缓存(开发/测试用)和Redis(生产多实例用)的统一方案,无需额外定制“伪实现”。
标准替代方案:基于「状态标记+分布式锁」的间接控制
不需要修改IDistributedCache的实现,而是通过分布式缓存存储取消状态标记,配合任务内部轮询标记实现取消逻辑,同时用分布式锁保证任务单实例执行。这套方案完全兼容AddDistributedMemoryCache和Redis的AddStackExchangeRedisCache,适配未来扩展需求。
具体步骤:
- 用
IDistributedCache存储两个关键状态:- 任务运行状态标记(比如
"LongRunningTask:IsRunning"):控制多实例下只有一个实例启动任务。 - 任务取消标记(比如
"LongRunningTask:CancellationRequested"):跨实例触发任务取消。
- 任务运行状态标记(比如
- 任务内部逻辑:
- 启动前通过分布式锁抢占执行权,只有抢到锁的实例才启动任务。
- 任务执行过程中定期轮询缓存中的取消标记,一旦标记存在,触发本地
CancellationTokenSource取消任务。
- 取消操作入口:
需要取消任务时,直接往分布式缓存写入取消标记(设置短过期时间避免残留),所有实例的任务都会读取标记并响应取消。
代码示例(兼容本地/Redis缓存)
1. 注册缓存服务(开发/生产无缝切换)
开发环境用内存缓存:
services.AddDistributedMemoryCache();
生产环境替换为Redis缓存:
services.AddStackExchangeRedisCache(options => { options.Configuration = "your-redis-connection-string"; options.InstanceName = "YourApp:"; });
2. 任务管理类实现
public class LongRunningTaskManager { private readonly IDistributedCache _cache; private readonly ILogger<LongRunningTaskManager> _logger; private CancellationTokenSource _cts; private bool _isRunningLocally; private const string TaskRunningKey = "LongRunningTask:IsRunning"; private const string TaskCancelKey = "LongRunningTask:CancellationRequested"; private const int LockExpirySeconds = 30; // 分布式锁过期时间,需大于任务轮询间隔 public LongRunningTaskManager(IDistributedCache cache, ILogger<LongRunningTaskManager> logger) { _cache = cache; _logger = logger; } public async Task StartTaskAsync() { // 尝试获取分布式锁,确保单实例执行 var lockValue = Guid.NewGuid().ToString(); var lockAcquired = await TryAcquireLockAsync(TaskRunningKey, lockValue, TimeSpan.FromSeconds(LockExpirySeconds)); if (!lockAcquired) { _logger.LogInformation("任务已在其他实例运行,当前实例不启动"); return; } try { _isRunningLocally = true; _cts = new CancellationTokenSource(); _logger.LogInformation("开始执行长时间运行任务"); // 模拟长时间任务,定期检查取消标记 while (!_cts.Token.IsCancellationRequested) { // 读取分布式缓存中的取消标记 var cancelRequested = await _cache.GetStringAsync(TaskCancelKey); if (!string.IsNullOrEmpty(cancelRequested)) { _cts.Cancel(); break; } // 执行实际任务逻辑 await DoWorkAsync(_cts.Token); // 刷新分布式锁,防止任务未完成时锁过期 await _cache.SetStringAsync(TaskRunningKey, lockValue, new DistributedCacheEntryOptions { AbsoluteExpirationRelativeToNow = TimeSpan.FromSeconds(LockExpirySeconds) }); await Task.Delay(1000, _cts.Token); // 轮询间隔 } } catch (OperationCanceledException) { _logger.LogInformation("任务已被取消"); } finally { _isRunningLocally = false; // 仅当前持有锁的实例清理状态 var currentLockValue = await _cache.GetStringAsync(TaskRunningKey); if (currentLockValue == lockValue) { await _cache.RemoveAsync(TaskRunningKey); await _cache.RemoveAsync(TaskCancelKey); } _cts?.Dispose(); } } public async Task CancelTaskAsync() { // 写入取消标记到分布式缓存 await _cache.SetStringAsync(TaskCancelKey, "true", new DistributedCacheEntryOptions { AbsoluteExpirationRelativeToNow = TimeSpan.FromMinutes(5) // 标记过期时间,避免残留 }); // 当前实例正在运行时,直接触发本地取消 if (_isRunningLocally && _cts != null) { _cts.Cancel(); } } private async Task<bool> TryAcquireLockAsync(string key, string value, TimeSpan expiry) { // 分布式锁基础实现:SETNX逻辑 var existingValue = await _cache.GetStringAsync(key); if (!string.IsNullOrEmpty(existingValue)) { return false; } await _cache.SetStringAsync(key, value, new DistributedCacheEntryOptions { AbsoluteExpirationRelativeToNow = expiry }); // 二次检查避免并发问题 var newValue = await _cache.GetStringAsync(key); return newValue == value; } private async Task DoWorkAsync(CancellationToken token) { // 替换为实际任务逻辑 _logger.LogInformation("执行任务工作单元"); await Task.Delay(2000, token); } }
3. 注册任务管理类为单例
services.AddSingleton<LongRunningTaskManager>();
方案适配性说明
- 兼容IDistributedCache全生态:切换内存缓存/Redis仅需修改注册代码,业务逻辑无需变动。
- 规避CancellationToken序列化问题:用字符串标记替代内存对象,完美解决分布式跨实例通信。
- 自带故障恢复:若执行任务的实例崩溃,分布式锁过期后其他实例可重新抢占锁启动任务(可根据业务调整是否自动重启)。
内容的提问来源于stack exchange,提问作者Benyom
相关产品推荐
相关产品推荐

