基于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
相关产品推荐
相关产品推荐

