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

SignalR重连时Azure消息会话锁定问题求助

Azure Service Bus Session 锁释放异常问题

使用React/Redux前端与ASP.NET Core服务器建立SignalR连接,连接成功后前端触发startTesting函数,该函数通过循环每4秒向前端推送一次更新,数据来自Azure Message Session(需保证FIFO),通过单例消息处理器获取数据后由SignalR Hub转发至前端。

正常流程运行无问题,但SignalR连接刷新(旧连接释放、新连接创建)时,会抛出错误:The requested session cannot be accepted. It may be locked by another receiver.,即使在finally块中调用了await receiver.CloseAsync()。

单例消息处理器代码

public async Task<string> ReceiveMessage(string imei, CancellationToken ct = default) 
{ 
    var receiver = await _sessionClient.AcceptMessageSessionAsync(imei);

    try
    {
        if (!_processingSessions.TryAdd(imei, true))
        {
            throw new Exception("Session is already being processed");
        }

        await Task.Delay(1000, ct);

        Message message = await receiver.ReceiveAsync(TimeSpan.FromMilliseconds(10));

        if (message == null) return null;

        _logger.LogInformation("Message received: {Message}", message.Body);

        await receiver.CompleteAsync(message.SystemProperties.LockToken);
        
        return Encoding.UTF8.GetString(message.Body);
    }
    catch (Exception e)
    {
        _logger.LogError("Service Bus Exception: {ExMessage}", e.Message);
        throw new Exception(e.Message);
    }
    finally
    {
        // _logger.LogInformation("Closing Message Session");
        _processingSessions.TryRemove(imei, out var result);
        await receiver.CloseAsync();
    }
}

已尝试的ServiceBusSessionReceiver实现

// 构造函数:
// _serviceBusClient = serviceBusClient.CreateClient("DeviceTestQueue");

// 消息处理器接收逻辑:
// var receiver = await _serviceBusClient.AcceptSessionAsync(
//     _configuration.GetSection("DEVICETESTMSGQUEUE").Value,
//     imei, cancellationToken: ct
// );

前端SignalR代码

import {
    HttpTransportType,
    HubConnectionBuilder,
    HubConnectionState,
    JsonHubProtocol,
    LogLevel
} from '@microsoft/signalr';
import * as types from "../actions/actionsTypes";

const baseUrl = process.env.SOCKET_URL + 'devicetests/ws';
const protocol = new JsonHubProtocol();

let connection = null;

export const initializeConnection = (imei) => async (dispatch) => {
    let url = `not relevant`;

    if (!connection) {
        console.log("Creating new DeviceTestSocket connection");

        connection = new HubConnectionBuilder().withUrl(
            url,
            {
                // accessTokenFactory: () => accessToken,
                transport: HttpTransportType.WebSockets | HttpTransportType.LongPolling,
                logMessageContent: true,
                logger: LogLevel.Debug,
            })
            .withHubProtocol(protocol)
            .build();
    }

    await connection.start().then(() => {
        console.log("Connected to DeviceTestSocket");

        startTesting(imei);
    });


    await connection.on("SendTests", (message) => {
        console.log("SendTests Triggered");
        dispatch({
            type: types.INITIALIZE_TESTS,
            payload: {
                deviceTests: message
            }
        })
    });

    await connection.on("UpdateTest", (message) => {
        console.log("UpdateTest Triggered");
        dispatch({
            type: types.UPDATE_TEST,
            payload: {
                deviceTests: message
            }
        })
    });
}

export const startTesting = async (imei) => {
    if (isConnectionOpen) {
        await connection.invoke("StartTesting", imei, new Date());
    }
}

export const closeConnection = () => async (dispatch) => {
    if (isConnectionOpen) {
        connection.stop().then(() => {
            console.log("Disconnected from DeviceTestSocket");
        });

        dispatch({
            type: types.FINISH_TESTS,
            payload: {
                deviceTests: []
            }
        });
    }
}

export const isConnectionOpen = () => connection && connection.state === HubConnectionState.Connected;

补充:调用ReceiveMessage的SignalR Hub代码

采用Rick Montalvo的解决方案后情况有所改善,但严格测试仍会出现该错误,以下是Hub相关代码:

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);
    private readonly CancellationTokenSource _cts;


    public DeviceTestHub(IDeviceTestHandler deviceTestHandler, IMemoryCache memoryCache,
        ILogger<DeviceTestHub> logger, IMessageProcessor messageProcessor)
    {
        _cache = memoryCache;
        _logger = logger;
        _messageProcessor = messageProcessor;
        _deviceTestHandler = deviceTestHandler;
        
        // 设置12分钟的取消令牌,超时则取消测试
        _cts = new CancellationTokenSource(TimeSpan.FromMinutes(12));
    }

    public override async Task OnConnectedAsync()
    {
        // 从请求头获取IMEI
        var context = Context.GetHttpContext();

        // 尝试从查询参数获取imei
        if (!context!.Request.Query.TryGetValue("imei", out var imei)) return;

        _logger.LogInformation("Connection established for device: {imei}", imei);

        await base.OnConnectedAsync();

        var savedTests = _cache.Get<List<SocketModel>>($"{MemoryCacheKeys.DeviceTest}-{imei}");

        if (savedTests != null)
        {
            // 发送保存的测试数据到前端
            await Clients.Caller.SendTests(savedTests);

            return;
        }

        // 获取测试列表
        var socketTests = await _deviceTestHandler.GetTests(imei);

        // 缓存该IMEI对应的测试列表,用于断开重连场景
        _cache.Set($"{MemoryCacheKeys.DeviceTest}-{imei}", socketTests, _cacheExpirationTime);


        // 发送测试列表到前端
        await Clients.Caller.SendTests(socketTests);

        // 清理可能存在的旧会话
        await _messageProcessor.Dispose(imei);
    }

    public async Task StartTesting(string imei, DateTime startDate)
    {
        try
        {
            // 测试开始前清理会话
            await _messageProcessor.Dispose(imei);

            // 获取测试列表
            var socketTests =
                _cache.Get<List<SocketModel>>($"{MemoryCacheKeys.DeviceTest}-{imei}")
                ?? await _deviceTestHandler.GetTests(imei);

            var messages = new List<RawMessageDTO>();

            while (!_cts.IsCancellationRequested)
            {
                _cts.Token.ThrowIfCancellationRequested();

                if (socketTests.Any(x => x.State == State.NotStarted))
                {
                    // 找到第一个未开始的测试并标记为加载中
                    var test = _deviceTestHandler.SetTestActive(socketTests, _cts.Token);

                    // 推送测试状态更新到前端
                    if (test != null) await Clients.Caller.UpdateTest(test);
                }
                
                // 从Message Session获取消息
                var message = await _messageProcessor.ReceiveMessage(imei, _cts.Token);

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

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

                    // 推送原始消息到前端
                    await Clients.Caller.SendRawMessage(JsonSerializer.Serialize(rawMsg));

                    messages.Add(rawMsg);
                }

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

                if (activeTest is null) return;

                var executedTest = await _deviceTestHandler.ExecuteTest(imei, activeTest, messages, startDate);

                if (executedTest != null)
                {
                    // 推送测试执行结果到前端
                    await Clients.Caller.UpdateTest(executedTest);
                }
            }
        }
        catch (OperationCanceledException e)
        {
            await Clients.Caller.SendMessage($"Test took longer than 12 minutes and is cancelled for imei: {imei}");
            await _messageProcessor.Dispose(imei);
            await OnDisconnectedAsync(e);
        }
        catch (Exception ex)
        {
            _logger.LogError("Error occurred during testing: {Message}", ex.Message);
            await _messageProcessor.Dispose(imei);
            await Clients.Caller.SendMessage(ex.Message);
        }
    }

    public override async Task OnDisconnectedAsync(Exception? exception)
    {
        _cts.Cancel();
        
        // 从请求头获取IMEI
        var context = Context.GetHttpContext();

        // 尝试从查询参数获取imei
        if (!context!.Request.Query.TryGetValue("imei", out var imei)) return;

        await _messageProcessor.Dispose(imei);

        if (exception != null)
            _logger.LogError("Disconnected with error: {Message}", exception.Message);

        await base.OnDisconnectedAsync(exception);
    }
}

恳请帮忙排查并解决该问题!

内容的提问来源于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.15 18:44:53