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

基于ASP.NET Core与MSSQL Server的实时数据变更通知实现

实现MSSQL数据变更触发SignalR实时页面更新

一、启用SQL Server Service Broker

为避免直接在触发器中调用外部服务导致的性能和事务问题,先给目标数据库启用Service Broker,用于异步处理变更事件:

ALTER DATABASE YourDatabaseName SET ENABLE_BROKER WITH ROLLBACK IMMEDIATE;

二、创建变更捕获的触发器与Service Broker对象

通过Service Broker的消息队列机制捕获表变更,以下是针对目标表(示例表名YourTargetTable)的完整SQL脚本:

-- 定义消息类型
CREATE MESSAGE TYPE [ChangeNotificationMessage] VALIDATION = WELL_FORMED_XML;

-- 创建契约
CREATE CONTRACT [ChangeNotificationContract] ([ChangeNotificationMessage] SENT BY INITIATOR);

-- 创建存储消息的队列
CREATE QUEUE [ChangeNotificationQueue];

-- 创建绑定队列的服务
CREATE SERVICE [ChangeNotificationService] ON QUEUE [ChangeNotificationQueue] ([ChangeNotificationContract]);

-- 创建触发器,捕获插入/删除事件
CREATE TRIGGER Trigger_YourTargetTable_Change
ON YourTargetTable
AFTER INSERT, DELETE
AS
BEGIN
    SET NOCOUNT ON;

    DECLARE @DialogHandle UNIQUEIDENTIFIER;
    DECLARE @ChangeType NVARCHAR(10);
    DECLARE @MessageBody XML;

    -- 判断变更类型
    IF EXISTS(SELECT * FROM INSERTED)
        SET @ChangeType = 'INSERT';
    ELSE IF EXISTS(SELECT * FROM DELETED)
        SET @ChangeType = 'DELETE';

    -- 构造消息内容(可根据需求添加更多字段,比如变更记录的ID)
    SET @MessageBody = (
        SELECT @ChangeType AS ChangeType,
               (SELECT Id FROM INSERTED FOR XML PATH(''), ROOT('Inserted')) AS InsertedRecords,
               (SELECT Id FROM DELETED FOR XML PATH(''), ROOT('DeletedRecords')) AS DeletedRecords
        FOR XML PATH(''), ROOT('ChangeEvent')
    );

    -- 发送消息到Service Broker服务
    BEGIN DIALOG CONVERSATION @DialogHandle
        FROM SERVICE [ChangeNotificationService]
        TO SERVICE 'ChangeNotificationService'
        ON CONTRACT [ChangeNotificationContract]
        WITH ENCRYPTION = OFF;

    SEND ON CONVERSATION @DialogHandle MESSAGE TYPE [ChangeNotificationMessage] (@MessageBody);

    END CONVERSATION @DialogHandle;
END
GO

三、编写.NET后台服务监听队列并调用SignalR Hub

创建一个后台托管服务,持续读取Service Broker队列中的变更消息,再通过SignalR Hub推送给前端:

1. 后台服务类

using Microsoft.AspNetCore.SignalR;
using Microsoft.Extensions.Hosting;
using Microsoft.Data.SqlClient;
using System.Xml.Linq;

public class DbChangeListener : BackgroundService
{
    private readonly IHubContext<DataUpdateHub> _hubContext;
    private readonly string _dbConnectionString;

    public DbChangeListener(IHubContext<DataUpdateHub> hubContext, IConfiguration config)
    {
        _hubContext = hubContext;
        _dbConnectionString = config.GetConnectionString("YourDbConnection");
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        while (!stoppingToken.IsCancellationRequested)
        {
            using var connection = new SqlConnection(_dbConnectionString);
            await connection.OpenAsync(stoppingToken);

            // 从队列接收一条消息
            var receiveCmd = new SqlCommand(@"
                RECEIVE TOP(1) 
                    conversation_handle, 
                    message_body, 
                    message_type_name
                FROM ChangeNotificationQueue", connection);

            using var reader = await receiveCmd.ExecuteReaderAsync(stoppingToken);
            if (reader.Read())
            {
                var convHandle = reader.GetGuid(0);
                var messageXml = reader.GetSqlXml(1).Value;
                var xmlDoc = XDocument.Parse(messageXml);
                
                // 解析变更类型
                var changeType = xmlDoc.Root.Element("ChangeType").Value;
                
                // 推送变更通知到所有前端客户端
                await _hubContext.Clients.All.SendAsync(
                    "OnDataChanged", 
                    new { Type = changeType }, 
                    stoppingToken
                );

                // 结束对话,清理队列
                var endCmd = new SqlCommand(@"
                    END CONVERSATION @ConvHandle", connection);
                endCmd.Parameters.AddWithValue("@ConvHandle", convHandle);
                await endCmd.ExecuteNonQueryAsync(stoppingToken);
            }

            await Task.Delay(800, stoppingToken); // 轮询间隔可根据业务调整
        }
    }
}

2. 注册服务与SignalR

在Program.cs中注册后台服务和SignalR:

var builder = WebApplication.CreateBuilder(args);

// 注册后台服务
builder.Services.AddHostedService<DbChangeListener>();
// 注册SignalR
builder.Services.AddSignalR();

var app = builder.Build();

// 配置SignalR端点
app.MapHub<DataUpdateHub>("/dataUpdateHub");

app.Run();

四、实现SignalR Hub类

创建基础的Hub类,用于前端连接和消息推送:

using Microsoft.AspNetCore.SignalR;

public class DataUpdateHub : Hub
{
    // 可添加客户端连接/断开的自定义逻辑,比如记录在线用户
}

五、前端页面连接SignalR接收更新

使用SignalR的JavaScript客户端连接Hub,实时接收变更通知并更新页面:

// 初始化SignalR连接
const connection = new signalR.HubConnectionBuilder()
    .withUrl("/dataUpdateHub")
    .build();

// 监听后端推送的变更事件
connection.on("OnDataChanged", function(changeInfo) {
    console.log(`数据库发生${changeInfo.Type}操作`);
    // 这里编写页面更新逻辑,比如重新加载表格数据、显示通知
    refreshTargetTable();
});

// 启动连接
async function startConnection() {
    try {
        await connection.start();
        console.log("SignalR连接成功");
    } catch (err) {
        console.error("连接失败,5秒后重试:", err);
        setTimeout(startConnection, 5000);
    }
}

startConnection();

// 示例:刷新目标表格数据
function refreshTargetTable() {
    // 调用接口获取最新数据并渲染到页面
    fetch("/api/yourTableData")
        .then(res => res.json())
        .then(data => {
            // 更新DOM逻辑
        });
}

关键注意事项

  • 确保SQL Server登录账号拥有Service Broker和队列的操作权限
  • 触发器中仅做消息发送,避免复杂逻辑影响数据库性能
  • 可根据业务需求扩展消息内容,比如传递变更记录的具体字段值
  • 轮询间隔可根据变更频率调整,或使用WAITFOR(RECEIVE...)实现更高效的消息监听

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 21:05:37