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

GraphQL订阅WebSocket异常断开1006,无法持续监听新增告警

解决GraphQL订阅断开及持续监听新增告警的问题

问题根源

你当前的订阅方法返回IQueryable<Alert>,这本质是一次性查询,执行后返回所有符合条件的现有告警,WebSocket连接完成数据推送后就会关闭,无法监听后续新增的告警。

解决方案

需要将订阅方法的返回类型改为IObservable<Alert>,同时实现数据库变更监听,确保新增符合条件的告警时能推送给客户端。以下是具体实现步骤:

1. 修改订阅方法的返回类型与逻辑

调整Subscription1类中的OnAlertAdded方法,改用IObservable<Alert>返回可观察序列,同时处理现有数据推送和新增数据监听:

public class Subscription1
{
    [Subscribe]
    [Topic("AlertAdded")]
    public IObservable<Alert> OnAlertAdded(
        [Service] ApplicationDbContext context,
        [Service] ILogger<Subscription1> logger,
        int clientId)
    {
        // 先推送已存在的符合条件的告警
        var existingAlerts = context.AlertClients
            .Where(ac => ac.ClientId == clientId)
            .Select(ac => ac.Alert)
            .ToList();

        // 创建可观察序列:先推送现有数据,再监听新增数据
        return Observable.Create<Alert>(observer =>
        {
            // 推送现有告警
            foreach (var alert in existingAlerts)
            {
                observer.OnNext(alert);
            }

            // 监听数据库中Alert的新增操作
            var token = context.ChangeTracker.StateChanged += (sender, e) =>
            {
                if (e.Entry.Entity is Alert newAlert && e.Entry.State == EntityState.Added)
                {
                    // 检查新增告警是否属于当前clientId
                    var isClientAlert = context.AlertClients
                        .Any(ac => ac.AlertId == newAlert.Id && ac.ClientId == clientId);

                    if (isClientAlert)
                    {
                        observer.OnNext(newAlert);
                    }
                }
            };

            // 订阅取消时清理资源,移除事件监听
            return Disposable.Create(() =>
            {
                context.ChangeTracker.StateChanged -= token;
                logger.LogInformation("订阅已取消,停止监听告警新增");
            });
        });
    }
}

2. 确保告警新增时触发推送

当服务器端新增告警并关联到目标clientId时,需确保调用ApplicationDbContext的SaveChanges方法,让ChangeTracker捕获到新增实体。示例新增代码:

public async Task AddAlertToClient(int clientId, Alert newAlert)
{
    using var context = new ApplicationDbContext();
    context.Alerts.Add(newAlert);
    await context.SaveChangesAsync();
    
    // 关联告警到指定客户端
    context.AlertClients.Add(new AlertClient { ClientId = clientId, AlertId = newAlert.Id });
    await context.SaveChangesAsync();
}

3. 关键注意事项

  • 使用IObservable:这是实现GraphQL持续推送的核心,代表一个可不断产生新数据的序列。
  • 变更监听:通过EF Core的ChangeTracker.StateChanged事件捕获新增的Alert实体,过滤后推送给对应客户端。
  • 资源清理:订阅取消时移除事件监听,避免内存泄漏。
  • Topic区分:如果需要按客户端隔离推送,可以将Topic设置为$"AlertAdded_{clientId}",确保每个客户端只收到自己的告警。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 09:27:43