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

SignalR中取消长运行循环失败,求排查解决方法

问题排查与解决方案

核心问题分析

  1. Hub实例瞬时性导致令牌不共享:SignalR默认每次调用Hub方法都会创建新的Hub实例,你在StartTesting中创建的_cts是当前实例的私有字段,而调用CancelTesting时是全新的Hub实例,两者的_cts完全不是同一个对象,因此Cancel()操作根本不会作用到正在运行的测试任务。
  2. CancellationToken使用不规范:方法参数传入了ct,但代码中全程直接访问字段_cts.Token,既增加耦合性,也放大了实例不一致带来的问题。
  3. 取消检查粒度不足:部分耗时操作(如ExecuteTest)未嵌入取消令牌检查,导致取消请求无法及时响应。
  4. 未处理客户端断开的自动取消:当前代码未在客户端断开连接时触发任务取消逻辑。

具体修复方案

1. 用缓存存储CancellationTokenSource

利用你已有的IMemoryCache,将每个测试任务的令牌源与IMEI+连接ID关联存储,确保StartTesting和CancelTesting能访问到同一个令牌对象。

2. 规范CancellationToken使用

统一使用方法传入的CancellationToken,避免直接访问Hub字段,降低耦合。

3. 细化取消检查

在所有可能阻塞或耗时的操作前后添加取消检查,确保取消请求能被立即响应。

4. 处理客户端断开的自动取消

重写OnDisconnectedAsync方法,清理当前连接对应的所有测试任务。


修改后的完整代码

public class DeviceTestHub : Hub<IDeviceTestClient> 
{
    private readonly ILogger<DeviceTestHub> _logger;
    private readonly IMemoryCache _cache;
    private readonly IMessageProcessor _messageProcessor;
    private readonly IDeviceTestHandler _deviceTestHandler;
    private readonly TimeSpan _cacheExpirationTime = TimeSpan.FromMinutes(15);

    public DeviceTestHub(IDeviceTestHandler deviceTestHandler, IMemoryCache memoryCache,
        ILogger<DeviceTestHub> logger, IMessageProcessor messageProcessor)
    {
        _cache = memoryCache;
        _logger = logger;
        _messageProcessor = messageProcessor;
        _deviceTestHandler = deviceTestHandler;
    }

    public async Task StartTesting(string imei, DateTime startDate)
    {
        var cts = new CancellationTokenSource();
        // 用IMEI+连接ID作为缓存键,避免同IMEI多任务冲突
        var cacheKey = $"{MemoryCacheKeys.DeviceTestCts}-{imei}-{Context.ConnectionId}";
        _cache.Set(cacheKey, cts, _cacheExpirationTime);

        try
        {
            await LongRunningTest(imei, startDate, cts.Token, cacheKey);
        }
        catch (OperationCanceledException e)
        {
            await Clients.Caller.SendMessage($"测试超时12分钟,已取消IMEI: {imei}");
            await _messageProcessor.Dispose(imei);
            await base.OnDisconnectedAsync(e);
        }
        catch (Exception ex)
        {
            _logger.LogError($"WebSocket异常,IMEI: {imei}, 异常信息: {ex.Message}");
            await _messageProcessor.Dispose(imei);
            await Clients.Group(imei).SendMessage(ex.Message);
        }
        finally
        {
            cts.Dispose();
            _cache.Remove(cacheKey);
        }
    }

    public void CancelTesting(string imei)
    {
        _logger.LogInformation("================================= 取消测试 =================================");
        var cacheKey = $"{MemoryCacheKeys.DeviceTestCts}-{imei}-{Context.ConnectionId}";
        if (_cache.TryGetValue<CancellationTokenSource>(cacheKey, out var cts))
        {
            cts.Cancel();
        }
    }

    public override async Task OnDisconnectedAsync(Exception exception)
    {
        // 清理当前连接对应的所有测试令牌
        var cacheKeys = _cache.GetKeys().Where(k => k.EndsWith(Context.ConnectionId));
        foreach (var key in cacheKeys)
        {
            if (_cache.TryGetValue<CancellationTokenSource>(key, out var cts))
            {
                cts.Cancel();
                cts.Dispose();
            }
            _cache.Remove(key);
        }
        await _messageProcessor.Dispose(Context.ConnectionId);
        await base.OnDisconnectedAsync(exception);
    }

    private async Task LongRunningTest(string imei, DateTime startDate, CancellationToken ct, string cacheKey)
    {
        await _messageProcessor.Dispose(imei);

        var socketTests =
            _cache.Get<List<SocketModel>>($"{MemoryCacheKeys.DeviceTest}-{imei}")
            ?? await _deviceTestHandler.GetTests(imei);

        var messages = new List<RawMessageDTO>();
        
        ct.ThrowIfCancellationRequested();

        while (!ct.IsCancellationRequested)
        {
            ct.ThrowIfCancellationRequested();

            if (socketTests.Any(x => x.State == State.NotStarted))
            {
                var test = _deviceTestHandler.SetTestActive(socketTests, ct);
                ct.ThrowIfCancellationRequested();

                if (test != null) await Clients.Group(imei).UpdateTest(test);
            }

            ct.ThrowIfCancellationRequested();
            var message = await _messageProcessor.ReceiveMessage(imei, ct);

            ct.ThrowIfCancellationRequested();
            if (message.IsNullOrEmpty() && messages.IsNullOrEmpty())
            {
                await Task.Delay(3000, ct);
                continue;
            }

            ct.ThrowIfCancellationRequested();
            if (!message.IsNullOrEmpty())
            {
                var rawMsg = await _deviceTestHandler.ConvertMessageToRawMessage(message);
                ct.ThrowIfCancellationRequested();

                await Clients.Caller.SendRawMessage(JsonSerializer.Serialize(rawMsg));
                messages.Add(rawMsg);
            }

            ct.ThrowIfCancellationRequested();
            var activeTest = socketTests.FirstOrDefault(x => x.State is State.Loading or State.Inprogress);

            if (activeTest is null) return;

            ct.ThrowIfCancellationRequested();
            // 注意:需要给ExecuteTest添加CancellationToken参数并内部检查取消
            var executedTest = await _deviceTestHandler.ExecuteTest(imei, activeTest, messages, startDate, ct);
            ct.ThrowIfCancellationRequested();

            if (executedTest != null)
            {
                await Clients.Group(imei).UpdateTest(executedTest);
            }
        }
    }
}

额外注意事项

  • ExecuteTest方法改造:必须给_deviceTestHandler.ExecuteTest添加CancellationToken参数,并在方法内部的关键节点(如循环、耗时IO操作)添加取消检查,否则长时间运行的测试逻辑无法响应取消。
  • 内存泄漏防护:任务完成、取消或客户端断开时,务必清理缓存中的CancellationTokenSource并释放资源。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 20:05:54