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
相关产品推荐
相关产品推荐

