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

多实例微服务:如何实现跨实例任务取消?

解决方案思路

核心问题拆解

你面临的核心矛盾是:分布式环境下要跨实例控制任务取消,但CancellationToken是内存内对象,没法直接序列化存入分布式缓存。而你需要一套能兼容本地内存缓存(开发/测试用)和Redis(生产多实例用)的统一方案,无需额外定制“伪实现”。

标准替代方案:基于「状态标记+分布式锁」的间接控制

不需要修改IDistributedCache的实现,而是通过分布式缓存存储取消状态标记,配合任务内部轮询标记实现取消逻辑,同时用分布式锁保证任务单实例执行。这套方案完全兼容AddDistributedMemoryCache和Redis的AddStackExchangeRedisCache,适配未来扩展需求。

具体步骤:

  • 用IDistributedCache存储两个关键状态:
    1. 任务运行状态标记(比如"LongRunningTask:IsRunning"):控制多实例下只有一个实例启动任务。
    2. 任务取消标记(比如"LongRunningTask:CancellationRequested"):跨实例触发任务取消。
  • 任务内部逻辑:
    1. 启动前通过分布式锁抢占执行权,只有抢到锁的实例才启动任务。
    2. 任务执行过程中定期轮询缓存中的取消标记,一旦标记存在,触发本地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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 22:52:18