SignalR客户端网络故障场景下消息投递的健壮处理方案咨询
基于.NET Core 7 + React Native SignalR的离线消息可靠投递方案
问题背景
我开发了一个基于.NET Core 7的后端项目,与React Native移动端应用通过SignalR通信。应用用户常在农村地区开展产品销售,频繁遇到网络短暂中断或后端发消息时断网的情况,导致消息丢失,需要解决两类场景下的消息投递问题:
服务端初始配置
using Microsoft.EntityFrameworkCore; using UbuntuLifeBackendCore.Hubs; using UbuntuLifeBackendCore.Models; var builder = WebApplication.CreateBuilder(args); builder.Services.AddControllers(); builder.Services.AddEndpointsApiExplorer(); builder.Services.AddSwaggerGen(); builder.Services.AddSignalR(hubOptions => { hubOptions.ClientTimeoutInterval = TimeSpan.FromSeconds(15); // 客户端到服务端超时 hubOptions.KeepAliveInterval = TimeSpan.FromSeconds(5); // 心跳包频率 hubOptions.EnableDetailedErrors = true; }).AddJsonProtocol(options => { options.PayloadSerializerOptions.PropertyNamingPolicy = null; }); builder.Services.AddSingleton<HubHelper>(); builder.Services.AddDbContext<ApplicationDbContext>(options => options.UseSqlServer(builder.Configuration.GetConnectionString("LiveConnection"))); var app = builder.Build(); if (app.Environment.IsDevelopment()) { app.UseSwagger(); app.UseSwaggerUI(); } app.UseHttpsRedirection(); app.UseAuthorization(); app.MapControllers(); app.UseCors(x => x .AllowAnyMethod() .AllowAnyHeader() .SetIsOriginAllowed(origin => true) .AllowCredentials()); app.MapHub<ProductHub>("/hubs/productsHub"); app.MapHub<ClientHub>("/hubs/clientsHub"); app.Run();
待解决的两类问题
1. 客户端完全断网重连后的消息补发
当前客户端连接/重连后会调用RegisterUser,将ConnectionId存入数据库,并尝试补发未投递消息,但现有逻辑不完整。
服务端RegisterUser现有代码
// Server side code for Register User function public async Task RegisterUser(string token) { try { // saves the data of a client in the database for their connection Id // method to pull the logs in the database for undelivered messages and sending them based on their use cases } catch (Exception ex) { // logging error } }
React Native客户端SignalR配置
export const initSignalR = async (signalRURL: string, authKey: string) => { // initialize signalR builder connection hub connectionHub = new signalR.HubConnectionBuilder() .configureLogging(signalR.LogLevel.Trace) .withUrl(signalRURL, signalR.HttpTransportType.WebSockets) .withAutomaticReconnect({ nextRetryDelayInMilliseconds: (retryContext) => { console.log(`Retry :: ${retryContext.previousRetryCount}`) return 2000 }, }) .build() const startConnection = () => { connectionHub .start() .then(() => { connectionSucceed(authKey) }) .catch((err) => { console.error("Connection failed: ", err) setTimeout(() => { startConnection() }, 2000) }) } // start connection startConnection() // reconnect connection connectionHub.onreconnected(() => { connectionSucceed(authKey) }) return connectionHub } const receiveMessageHandler = (messageId: any, onSuccess?: any) => { connectionHub .invoke("ReceivedMessage", messageId) .then((res) => { onSuccess?.(res) }) .catch((err) => { if (connectionHub?.state == "Connected") { receiveMessageHandler(messageId, onSuccess) } }) } const connectionSucceed = (authKey: string) => { connectionHub .invoke("RegisterUser", authKey) .then(() => { }) .catch((err) => { console.error("Error registering user after reconnection:", err) }) }
客户端消息处理代码
messageOnSignalR("ReceiveUploadVoiceMandateResponse", (res1) => { let res = JSON.parse(res1) console.log("Recieved response from server on => ReceiveUploadVoiceMandateResponse", res) try { receiveMessageFromSignalR(res?.messageId, (receiveMessageRes) => { if (res.Status) { console.log("Voice Uploaded successfully") } }) } catch (e) {} })
2. 客户端短暂断网时的消息丢失
曾尝试用Windows服务每5秒检查补发消息,但无法满足实时性需求,需要无需Windows服务的高效健壮方案。
解决方案
一、完全断网重连的消息补发优化
1. 完善消息持久化模型
首先在数据库中创建消息表,记录关键信息:
public class Message { public int Id { get; set; } public string UserId { get; set; } // 关联用户ID public string MethodName { get; set; } // 客户端要调用的SignalR方法名 public string Payload { get; set; } // 消息内容(JSON格式) public MessageStatus Status { get; set; } // 待投递/已投递/已过期 public DateTime CreatedAt { get; set; } = DateTime.UtcNow; } public enum MessageStatus { Pending, Delivered, Expired } public class UserConnection { public string UserId { get; set; } public string ConnectionId { get; set; } }
2. 补全RegisterUser逻辑,实现重连补发
private readonly ApplicationDbContext _context; private readonly ILogger<ClientHub> _logger; // 通过构造函数注入依赖 public ClientHub(ApplicationDbContext context, ILogger<ClientHub> logger) { _context = context; _logger = logger; } public async Task RegisterUser(string token) { try { // 解析token获取用户ID(替换为你的实际认证逻辑) var userId = JwtTokenParser.ParseUserId(token); // 更新用户当前ConnectionId var userConnection = await _context.UserConnections .FirstOrDefaultAsync(u => u.UserId == userId); if (userConnection == null) { _context.UserConnections.Add(new UserConnection { UserId = userId, ConnectionId = Context.ConnectionId }); } else { userConnection.ConnectionId = Context.ConnectionId; } await _context.SaveChangesAsync(); // 查询24小时内未投递的消息,按时间顺序补发 var undeliveredMessages = await _context.Messages .Where(m => m.UserId == userId && m.Status == MessageStatus.Pending && m.CreatedAt > DateTime.UtcNow.AddHours(-24)) .OrderBy(m => m.CreatedAt) .ToListAsync(); foreach (var message in undeliveredMessages) { try { // 调用客户端对应的方法发送消息 await Clients.Client(Context.ConnectionId) .SendAsync(message.MethodName, message.Payload); message.Status = MessageStatus.Delivered; } catch (Exception ex) { _logger.LogError(ex, "补发消息失败:MessageId={0}, UserId={1}", message.Id, userId); // 单个消息失败不中断后续补发 } } await _context.SaveChangesAsync(); } catch (Exception ex) { _logger.LogError(ex, "用户注册及消息补发失败"); } }
3. 强化客户端确认机制
客户端收到消息后调用ReceivedMessage,服务端仅在收到确认后更新状态,避免重复补发:
public async Task ReceivedMessage(int messageId) { var message = await _context.Messages.FindAsync(messageId); if (message != null && message.Status == MessageStatus.Pending) { message.Status = MessageStatus.Delivered; await _context.SaveChangesAsync(); } }
二、短暂断网的消息可靠投递
1. 优化SignalR连接配置
调整服务端超时参数,适配农村不稳定网络:
builder.Services.AddSignalR(hubOptions => { hubOptions.ClientTimeoutInterval = TimeSpan.FromSeconds(30); // 延长客户端超时 hubOptions.KeepAliveInterval = TimeSpan.FromSeconds(10); // 降低心跳频率,减少网络消耗 hubOptions.EnableDetailedErrors = true; })
客户端优化重连策略,使用指数退避避免频繁重试:
.withAutomaticReconnect({ nextRetryDelayInMilliseconds: (retryContext) => { const delays = [1000, 2000, 4000, 8000]; const delay = delays[retryContext.previousRetryCount]; console.log(`第${retryContext.previousRetryCount+1}次重连,延迟${delay}ms`); return delay || 8000; }, })
2. 服务端内存缓冲队列
为每个用户维护内存队列,连接不稳定时暂存消息,重连后自动发送:
// 静态字典存储用户消息队列,键为用户ID private static readonly ConcurrentDictionary<string, Queue<Message>> _userMessageQueues = new(); public async Task SendMessageToUser(string userId, string methodName, string payload) { var message = new Message { UserId = userId, MethodName = methodName, Payload = payload, Status = MessageStatus.Pending }; // 先持久化到数据库,保证消息不丢失 _context.Messages.Add(message); await _context.SaveChangesAsync(); // 获取用户当前连接ID var userConnection = await _context.UserConnections .FirstOrDefaultAsync(u => u.UserId == userId); if (userConnection != null) { var connectionId = userConnection.ConnectionId; try { // 尝试直接发送 await Clients.Client(connectionId).SendAsync(methodName, payload); message.Status = MessageStatus.Delivered; await _context.SaveChangesAsync(); } catch { // 发送失败,加入内存队列 _userMessageQueues.AddOrUpdate(userId, queue => new Queue<Message>(new[] { message }), (_, queue) => { queue.Enqueue(message); return queue; }); } } } // 在RegisterUser中添加队列消息发送逻辑 // 补发完数据库消息后,检查内存队列 if (_userMessageQueues.TryGetValue(userId, out var queue)) { while (queue.Count > 0) { var queuedMessage = queue.Dequeue(); try { await Clients.Client(Context.ConnectionId).SendAsync(queuedMessage.MethodName, queuedMessage.Payload); queuedMessage.Status = MessageStatus.Delivered; _context.Messages.Update(queuedMessage); } catch (Exception ex) { _logger.LogError(ex, "发送队列消息失败:MessageId={0}", queuedMessage.Id); queue.Enqueue(queuedMessage); // 放回队列,下次重连再试 break; } } await _context.SaveChangesAsync(); }
3. 客户端本地缓存
客户端收到消息后先缓存到本地(AsyncStorage),处理完成后再通知服务端确认,避免APP重启丢失未处理消息:
messageOnSignalR("ReceiveUploadVoiceMandateResponse", async (res1) => { let res = JSON.parse(res1); console.log("Recieved response from server on => ReceiveUploadVoiceMandateResponse", res); // 先缓存到本地 await AsyncStorage.setItem(`signalr_message_${res.messageId}`, JSON.stringify(res)); try { receiveMessageFromSignalR(res?.messageId, async (receiveMessageRes) => { if (res.Status) { console.log("Voice Uploaded successfully"); // 处理完成后删除本地缓存 await AsyncStorage.removeItem(`signalr_message_${res.messageId}`); } }); } catch (e) {} }); // 客户端启动/重连时,检查本地未处理的消息 const checkLocalMessages = async () => { const keys = await AsyncStorage.getAllKeys(); const messageKeys = keys.filter(key => key.startsWith("signalr_message_")); for (const key of messageKeys) { const messageStr = await AsyncStorage.getItem(key); if (messageStr) { const message = JSON.parse(messageStr); // 重新触发处理逻辑 receiveMessageFromSignalR(message.messageId, () => { AsyncStorage.removeItem(key); }); } } }; // 在connectionSucceed中调用 const connectionSucceed = async (authKey: string) => { await checkLocalMessages(); connectionHub .invoke("RegisterUser", authKey) .then(() => {}) .catch((err) => { console.error("Error registering user after reconnection:", err); }); };
内容的提问来源于stack exchange,提问作者mihir doshi
相关产品推荐
相关产品推荐

