如何在Azure SQL Server特定表更新时触发EF Core服务执行
实现方案:Azure SQL表操作触发Service Bus消息(EF Core + 数据库触发器同步)
一、核心设计思路
采用「EF Core SaveChanges拦截器 + Azure SQL AFTER触发器」的组合方案,覆盖EF Core发起操作和直接数据库操作两种场景,确保数据完全写入数据库后再发送Service Bus消息,避免数据回滚但消息已发送的一致性问题。
二、EF Core层面实现(针对EF发起的增删改)
1. 自定义SaveChanges拦截器
创建继承自SaveChangesInterceptor的拦截器,重写事务提交后的SavedChanges/SavedChangesAsync方法,筛选目标表的变更并发送消息:
public class TargetTableChangeInterceptor : SaveChangesInterceptor { private readonly ServiceBusSender _serviceBusSender; public TargetTableChangeInterceptor(ServiceBusSender serviceBusSender) { _serviceBusSender = serviceBusSender; } public override int SavedChanges(SaveChangesCompletedEventData eventData) { var result = base.SavedChanges(eventData); ProcessTableChanges(eventData.Context); return result; } public override ValueTask<int> SavedChangesAsync(SaveChangesCompletedEventData eventData, CancellationToken cancellationToken = default) { var result = base.SavedChangesAsync(eventData, cancellationToken); ProcessTableChanges(eventData.Context); return result; } private void ProcessTableChanges(DbContext context) { // 筛选目标实体的新增/删除/更新记录 var added = context.ChangeTracker.Entries<TargetEntity>() .Where(e => e.State == EntityState.Added) .Select(e => e.Entity); var deleted = context.ChangeTracker.Entries<TargetEntity>() .Where(e => e.State == EntityState.Deleted) .Select(e => e.Entity); var updated = context.ChangeTracker.Entries<TargetEntity>() .Where(e => e.State == EntityState.Modified) .Select(e => new { Original = e.OriginalValues.ToObject<TargetEntity>(), Current = e.Entity }); // 发送对应操作的消息 if (added.Any()) { _serviceBusSender.SendMessageAsync( new ServiceBusMessage(JsonSerializer.Serialize(new { Operation = "Add", Data = added })) ).Wait(); } if (deleted.Any()) { _serviceBusSender.SendMessageAsync( new ServiceBusMessage(JsonSerializer.Serialize(new { Operation = "Delete", Data = deleted })) ).Wait(); } if (updated.Any()) { _serviceBusSender.SendMessageAsync( new ServiceBusMessage(JsonSerializer.Serialize(new { Operation = "Update", Data = updated })) ).Wait(); } } }
2. 注册拦截器与Service Bus客户端
在Program.cs中注入依赖:
builder.Services.AddDbContext<AppDbContext>(options => { options.UseSqlServer(builder.Configuration.GetConnectionString("AzureSql")) .AddInterceptors(builder.Services.BuildServiceProvider().GetRequiredService<TargetTableChangeInterceptor>()); }); // 注册Service Bus发送者 builder.Services.AddSingleton<ServiceBusSender>(sp => { var sbConn = builder.Configuration.GetConnectionString("ServiceBus"); return new ServiceBusClient(sbConn).CreateSender("table-change-topic"); }); builder.Services.AddSingleton<TargetTableChangeInterceptor>();
三、Azure SQL触发器实现(针对直接数据库操作)
创建AFTER类型触发器,确保数据写入数据库后触发消息发送:
1. 先创建Service Bus消息发送存储过程
CREATE PROCEDURE dbo.SendTableChangeMessage @operation NVARCHAR(10), @data NVARCHAR(MAX) AS BEGIN SET NOCOUNT ON; DECLARE @sbEndpoint NVARCHAR(200) = 'https://your-sb-namespace.servicebus.windows.net/table-change-topic/messages'; DECLARE @authHeader NVARCHAR(MAX) = '{"Authorization":"SharedAccessSignature sr=your-sb-namespace.servicebus.windows.net&sig=your-sas-key&se=1720000000&skn=RootManageSharedAccessKey", "Content-Type":"application/json"}'; DECLARE @body NVARCHAR(MAX) = JSON_QUERY('{"Operation":"' + @operation + '", "Data":' + @data + '}'); -- 调用Service Bus REST API发送消息 EXEC sp_invoke_external_rest_endpoint @url = @sbEndpoint, @method = 'POST', @headers = @authHeader, @body = @body; END
2. 创建增删改对应的AFTER触发器
-- 新增触发器 CREATE TRIGGER trg_TargetTable_AfterInsert ON dbo.TargetTable AFTER INSERT AS BEGIN SET NOCOUNT ON; DECLARE @data NVARCHAR(MAX) = (SELECT * FROM inserted FOR JSON AUTO); EXEC dbo.SendTableChangeMessage 'Add', @data; END; -- 删除触发器 CREATE TRIGGER trg_TargetTable_AfterDelete ON dbo.TargetTable AFTER DELETE AS BEGIN SET NOCOUNT ON; DECLARE @data NVARCHAR(MAX) = (SELECT * FROM deleted FOR JSON AUTO); EXEC dbo.SendTableChangeMessage 'Delete', @data; END; -- 更新触发器 CREATE TRIGGER trg_TargetTable_AfterUpdate ON dbo.TargetTable AFTER UPDATE AS BEGIN SET NOCOUNT ON; DECLARE @data NVARCHAR(MAX) = (SELECT i.*, d.* AS Original FROM inserted i JOIN deleted d ON i.Id = d.Id FOR JSON AUTO); EXEC dbo.SendTableChangeMessage 'Update', @data; END;
四、同步触发与一致性保障
- EF Core拦截器的
SavedChanges方法在数据库事务提交后执行,确保数据已持久化再发消息 - Azure SQL的AFTER触发器在数据写入(事务提交完成)后触发,避免未提交操作触发消息
- 若需严格同步,可使用示例中的同步消息发送方式;若允许异步,可改用后台任务,但需开启Service Bus消息持久化避免丢失
五、注意事项
- Service Bus的SAS密钥建议用Azure Key Vault管理,在存储过程中通过
AZURE_KEY_VAULT函数获取 - 触发器和拦截器中需添加异常捕获逻辑,避免因消息发送失败导致数据库操作回滚
- 高并发场景下,可考虑批量发送消息减少Service Bus调用次数
内容的提问来源于stack exchange,提问作者Vorasoda
相关产品推荐
相关产品推荐

