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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 13:12:55