ASP.NET Core SignalR:队列传组名至后台服务仅收一条消息
问题描述
我有一个使用SignalR的ASP.NET Core项目,尝试通过队列将组名作为参数传递给BackgroundService。在Hub中编写了从JavaScript传入组名加入群组的方法,但浏览器仅收到一条消息。
相关代码
LiveDataHub.cs
public sealed class LiveDataHub : Hub { private readonly IMyBackgroundQueue _queue; public LiveDataHub(IMyBackgroundQueue queue) => _queue = queue; public async Task JoinToGroup(string group) { Console.WriteLine("New group joined: " + group); await _queue.QueueAsync(group); await Groups.AddToGroupAsync(Context.ConnectionId, group); await Clients.Group(group).SendAsync("Send", $"{Context.ConnectionId} has joined the group {group}."); } }
MyBackgroundQueue.cs
public interface IMyBackgroundQueue { ValueTask QueueAsync(string item); IAsyncEnumerable<string> DequeueAllAsync( CancellationToken cancellationToken); } public class MyBackgroundQueue : IMyBackgroundQueue { //private readonly Channel<string> _channel = Channel.CreateUnbounded<string>(); private readonly Channel<string> _channel; public MyBackgroundQueue(int capacity) { var options = new BoundedChannelOptions(capacity) { FullMode = BoundedChannelFullMode.Wait }; _channel = Channel.CreateBounded<string>(options); } public ValueTask QueueAsync(string item) => _channel.Writer.WriteAsync(item); public IAsyncEnumerable<string> DequeueAllAsync(CancellationToken ct) => _channel.Reader.ReadAllAsync(ct); }
MyBackgroundService.cs
public class MyBackgroundService : BackgroundService { private readonly IHubContext<LiveDataHub> _hub; public IMyBackgroundQueue _queue { get; } public MyBackgroundService(IHubContext<LiveDataHub> hub, IMyBackgroundQueue queue) { _hub = hub; _queue = queue; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { while (!stoppingToken.IsCancellationRequested) { var eventMessage = new Models.EventMessage($"Id_{ Guid.NewGuid():N}", $"Title_{Guid.NewGuid():N}", DateTime.Now.ToString("yyyy-MM-dd hh:mm:ss")); await foreach(var _item in _queue.DequeueAllAsync(stoppingToken)) { Console.WriteLine("group name: " + _item.ToString()); await _hub.Clients.Group(_item.ToString()).SendAsync("onMessageReceived", eventMessage, stoppingToken); } await Task.Delay(1000, stoppingToken); } } public override async Task StopAsync(CancellationToken stoppingToken) { await base.StopAsync(stoppingToken); } }
Startup.cs
public void ConfigureServices(IServiceCollection services) { services.AddRazorPages(); services.AddSignalR(hubPotions=> { hubPotions.EnableDetailedErrors = true; }); services.AddHostedService<MyBackgroundService>(); services.AddSingleton<IMyBackgroundQueue, MyBackgroundQueue>(ctx => { return new MyBackgroundQueue(100); }); // 我怀疑错误在此行,但无法确定原因 }
javascript.js
const signalrConnection = new signalR.HubConnectionBuilder() .withUrl("/messagebroker") .configureLogging(signalR.LogLevel.Information) .build(); signalrConnection.start().then(function () { console.log("SignalR Hub Connected"); signalrConnection.invoke("JoinToGroup", "Group01") .catch(function (err) { console.error(err.toString()); }); }).catch(function (err) { signalrConnection.stop(); console.error(err.toString()); });
问题原因分析
核心问题出在MyBackgroundService的ExecuteAsync方法中:
await foreach(var _item in _queue.DequeueAllAsync(stoppingToken))会持续枚举通道内的元素,直到通道的Writer被标记为完成(调用Complete)。但你的场景中通道是持续接收组名的,这会导致程序第一次进入该循环后就永久阻塞,再也无法执行后续的Task.Delay和新消息生成逻辑。- 你仅在循环外创建一条
eventMessage,即使队列能取出多个组名,也只会发送同一条消息,且后续定时逻辑完全无法执行。
修复方案
方案1:调整后台服务循环逻辑
修改MyBackgroundService.cs,将队列处理与定时消息生成分离,避免阻塞:
public class MyBackgroundService : BackgroundService { private readonly IHubContext<LiveDataHub> _hub; public IMyBackgroundQueue _queue { get; } private readonly HashSet<string> _subscribedGroups = new HashSet<string>(); private readonly object _lock = new object(); public MyBackgroundService(IHubContext<LiveDataHub> hub, IMyBackgroundQueue queue) { _hub = hub; _queue = queue; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { // 启动独立任务处理队列中的组名订阅 var queueProcessingTask = ProcessQueue(stoppingToken); while (!stoppingToken.IsCancellationRequested) { var eventMessage = new Models.EventMessage($"Id_{ Guid.NewGuid():N}", $"Title_{Guid.NewGuid():N}", DateTime.Now.ToString("yyyy-MM-dd HH:mm:ss")); // 给所有已订阅的组发送消息 lock (_lock) { foreach (var group in _subscribedGroups) { await _hub.Clients.Group(group).SendAsync("onMessageReceived", eventMessage, stoppingToken); } } await Task.Delay(1000, stoppingToken); } await queueProcessingTask; } private async Task ProcessQueue(CancellationToken stoppingToken) { await foreach (var group in _queue.DequeueAllAsync(stoppingToken)) { lock (_lock) { if (!_subscribedGroups.Contains(group)) { _subscribedGroups.Add(group); Console.WriteLine($"Group added: {group}"); } } } } public override async Task StopAsync(CancellationToken stoppingToken) { await base.StopAsync(stoppingToken); } }
方案2:直接维护订阅组集合(更简洁)
放弃队列传递组名,直接通过一个单例服务维护已订阅的组:
1. 新增组订阅服务
public interface IGroupSubscriptionService { void AddGroup(string groupName); IEnumerable<string> GetAllGroups(); } public class GroupSubscriptionService : IGroupSubscriptionService { private readonly HashSet<string> _groups = new HashSet<string>(); private readonly object _lock = new object(); public void AddGroup(string groupName) { lock (_lock) _groups.Add(groupName); } public IEnumerable<string> GetAllGroups() { lock (_lock) return _groups.ToList(); } }
2. 修改LiveDataHub.cs
public sealed class LiveDataHub : Hub { private readonly IGroupSubscriptionService _subscriptionService; public LiveDataHub(IGroupSubscriptionService subscriptionService) => _subscriptionService = subscriptionService; public async Task JoinToGroup(string group) { Console.WriteLine("New group joined: " + group); await Groups.AddToGroupAsync(Context.ConnectionId, group); _subscriptionService.AddGroup(group); await Clients.Group(group).SendAsync("Send", $"{Context.ConnectionId} has joined the group {group}."); } }
3. 修改Startup.cs注册服务
public void ConfigureServices(IServiceCollection services) { services.AddRazorPages(); services.AddSignalR(hubPotions=> { hubPotions.EnableDetailedErrors = true; }); services.AddHostedService<MyBackgroundService>(); services.AddSingleton<IGroupSubscriptionService, GroupSubscriptionService>(); // 若不再需要队列,可移除以下注册 // services.AddSingleton<IMyBackgroundQueue, MyBackgroundQueue>(ctx => new MyBackgroundQueue(100)); }
4. 修改MyBackgroundService.cs
public class MyBackgroundService : BackgroundService { private readonly IHubContext<LiveDataHub> _hub; private readonly IGroupSubscriptionService _subscriptionService; public MyBackgroundService(IHubContext<LiveDataHub> hub, IGroupSubscriptionService subscriptionService) { _hub = hub; _subscriptionService = subscriptionService; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { while (!stoppingToken.IsCancellationRequested) { var eventMessage = new Models.EventMessage($"Id_{ Guid.NewGuid():N}", $"Title_{Guid.NewGuid():N}", DateTime.Now.ToString("yyyy-MM-dd HH:mm:ss")); foreach(var group in _subscriptionService.GetAllGroups()) { await _hub.Clients.Group(group).SendAsync("onMessageReceived", eventMessage, stoppingToken); } await Task.Delay(1000, stoppingToken); } } public override async Task StopAsync(CancellationToken stoppingToken) { await base.StopAsync(stoppingToken); } }
前端代码补充(必须添加消息监听)
你的前端缺少onMessageReceived的监听逻辑,添加后才能接收后台推送的消息:
const signalrConnection = new signalR.HubConnectionBuilder() .withUrl("/messagebroker") .configureLogging(signalR.LogLevel.Information) .build(); // 监听后台推送的实时消息 signalrConnection.on("onMessageReceived", function(message) { console.log("Received message:", message); // 此处可添加UI渲染逻辑 }); // 监听群组加入通知 signalrConnection.on("Send", function(message) { console.log(message); }); signalrConnection.start().then(function () { console.log("SignalR Hub Connected"); signalrConnection.invoke("JoinToGroup", "Group01") .catch(function (err) { console.error(err.toString()); }); }).catch(function (err) { signalrConnection.stop(); console.error(err.toString()); });
总结
- 原代码的核心问题是后台服务的循环被通道枚举操作阻塞,导致定时任务无法重复执行。
- 推荐使用直接维护订阅组集合的方案,更符合SignalR的组管理场景,避免队列带来的不必要复杂度。
- 前端必须添加对应的消息监听方法,才能接收到后台发送的消息。
内容的提问来源于stack exchange,提问作者SoftwareDveloper
相关产品推荐
相关产品推荐

