ASP.NET SignalR Clients.All.ReceiveMessage无法向所有客户端同步发消息
我实现了一个SignalR示例,功能是从第三方API获取最新更新,当API返回的LastUpdateDateTime属性时间戳发生变化时,控制器会向所有连接的客户端发送消息。控制器代码如下:
[Route("[controller]")] public class ChatController : Controller { private readonly IHubContext<ChatHub, IChatClient> _chatHub; List<Users> lstUsers = new List<Users>(); public ChatController(IHubContext<ChatHub, IChatClient> chatHub) { _chatHub = chatHub; } public async Task<IActionResult> Index() { return new ContentResult { Content = "Index Called", ContentType = "text/html" }; } [HttpGet("StartSocket")] public ContentResult StartSocket() { System.Timers.Timer timer = new System.Timers.Timer(1000); timer.Elapsed += (sender, e) => { timer.Stop(); timer.Dispose(); AppGlobal.SocketStarted = true; StartApiCalling(); }; timer.Start(); return new ContentResult { Content = "Socket Started", ContentType = "text/html" }; } [HttpGet("StopSocket")] public ContentResult StopSocket() { AppGlobal.SocketStarted = false; return new ContentResult { Content = "Socket Stopped", ContentType = "text/html" }; } private async Task CallApiAsync(string apiUrl) { Thread.Sleep(2000); while (AppGlobal.SocketStarted) { var client = new HttpClient(); var response = await client.GetAsync(apiUrl); var content = await response.Content.ReadAsStringAsync(); ProcessResponse(content); Thread.Sleep(50); } } private async Task StartApiCalling() { var tasks = new List<Task>(); var urls = new List<string> { "http://192.168.1.3/temp/api1", "http://192.168.1.3/temp/api2", "http://192.168.1.3/temp/api3" }; //reponse of api1 is like : [{"Id":"1","Name":"Name1","LastUpdateDateTime":"2023-05-07 17:28:42.334618"}] //reponse of api2 is like : [{"Id":"2","Name":"Name2","LastUpdateDateTime":"2023-05-07 17:28:42.334618"}] //reponse of api3 is like : [{"Id":"3","Name":"Name3","LastUpdateDateTime":"2023-05-07 17:28:42.334618"}] foreach (var url in urls) { tasks.Add(Task.Run(async () => { try { CallApiAsync(url); } catch (Exception ex) { Console.WriteLine(ex.Message); } })); } } private void ProcessResponse(string response) { if (string.IsNullOrEmpty(response)) { return; } List<Users> lstData = new List<Users>(); try { lstData = JsonConvert.DeserializeObject<List<Users>>(response) ?? new List<Users>(); } catch { } if (lstData != null && lstData.Count > 0) { foreach (var item in lstData) { var objUser = lstUsers.Where(x => x.Id == item.Id).FirstOrDefault(); if (objUser == null) { objUser = new Users(); objUser.DeepCopy(item); objUser.PropertyChanged += UserData_PropertyChanged; lstUsers.Add(objUser); } else { objUser.DeepCopy(item); } } } } private void UserData_PropertyChanged(object? sender, PropertyChangedEventArgs e) { if (sender != null) { Users usrs = (Users)sender; if (e.PropertyName == nameof(usrs.LastUpdateDateTime)) { SendUpdatedMessages(usrs.Name + " Changed - " + usrs.LastUpdateDateTime.ToString("dd/MMM/yyyy HH:mm:ss.fff")); } } } private async Task SendUpdatedMessages(string msg) { ChatMessage message = new ChatMessage { Group = "RoomA", Message = msg, User = "user1" }; await _chatHub.Clients.All.ReceiveMessage(message); }
问题
当有2个客户端连接时,两者收到的消息数量不同。例如运行约1分钟后,一个客户端收到279条消息,另一个收到254条,直到调用StopSocket操作。预期所有客户端收到的消息数量基本一致(若无法完全相同也应接近),因为每次API调用的时间戳都会变化,若消息丢失会导致客户端无法保持同步更新。请指出我遗漏或错误的地方。
异步任务未正确等待,导致任务脱离跟踪
在StartApiCalling方法中,调用CallApiAsync(url)时未使用await,外层Task.Run的异步lambda也未等待该任务,会导致CallApiAsync的异步操作脱离系统跟踪,可能出现任务被意外终止、调度异常的情况,进而引发消息推送丢失。修复时需添加await:try { await CallApiAsync(url); }非线程安全集合引发数据不一致
lstUsers是普通List<Users>,多个后台任务会同时调用ProcessResponse修改该集合,而List<T>不支持并发读写,会导致元素查找失败、更新丢失等问题,最终造成消息推送缺失。解决方式有两种:- 使用线程安全集合
ConcurrentDictionary<string, Users>(以用户Id为键)替代List<Users>; - 在操作
lstUsers时添加锁:private readonly object _usersLock = new object(); // ProcessResponse中操作集合的代码包裹在锁内 lock(_usersLock) { var objUser = lstUsers.Where(x => x.Id == item.Id).FirstOrDefault(); // 后续添加/更新逻辑 }
- 使用线程安全集合
HttpClient实例重复创建耗尽资源
在CallApiAsync的循环中每次创建新HttpClient实例,会快速耗尽套接字资源,导致API请求失败,进而丢失消息。应复用HttpClient实例,比如通过控制器注入IHttpClientFactory来创建实例,或使用静态HttpClient。异步推送任务未等待
UserData_PropertyChanged事件处理方法中调用SendUpdatedMessages(msg)时未使用await,会导致推送任务脱离上下文,可能在未完成推送时被中断,造成消息丢失。需将事件处理方法改为异步并等待推送:private async void UserData_PropertyChanged(object? sender, PropertyChangedEventArgs e) { if (sender != null) { Users usrs = (Users)sender; if (e.PropertyName == nameof(usrs.LastUpdateDateTime)) { await SendUpdatedMessages(usrs.Name + " Changed - " + usrs.LastUpdateDateTime.ToString("dd/MMM/yyyy HH:mm:ss.fff")); } } }Thread.Sleep阻塞线程池线程
异步方法中使用Thread.Sleep会阻塞线程池线程,影响任务调度效率,可能导致API请求和消息推送延迟或丢失。应替换为await Task.Delay(),比如Thread.Sleep(2000)改为await Task.Delay(2000),Thread.Sleep(50)改为await Task.Delay(50)。
内容的提问来源于stack exchange,提问作者Rohit

