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

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());
});

总结

  1. 原代码的核心问题是后台服务的循环被通道枚举操作阻塞,导致定时任务无法重复执行。
  2. 推荐使用直接维护订阅组集合的方案,更符合SignalR的组管理场景,避免队列带来的不必要复杂度。
  3. 前端必须添加对应的消息监听方法,才能接收到后台发送的消息。

内容的提问来源于stack exchange,提问作者SoftwareDveloper

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 23:17:28