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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 13:50:56