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

如何用C#/Azure函数实现数据库变更触发Webhook推送数据?

方案可行性分析

定时器轮询方案完全可行,但需根据你的实时性需求权衡:

  • 优势:实现简单,无需额外数据库配置,适合对延迟容忍度较高(如分钟级)的场景。
  • 劣势:存在轮询间隔内的延迟,若轮询过密会增加数据库负载。

如果需要近乎实时(秒级)的通知,更推荐变更数据捕获(CDC)+ 事件驱动的推模式:开启数据库CDC捕获变更事件,通过Azure Event Grid或Service Bus触发Azure函数,直接处理并推送Webhook,延迟更低且资源利用率更高。

Webhook核心实现逻辑

Webhook本质就是向客户指定的URL发起HTTP POST请求,将结构化的JSON数据作为请求体发送。关键要点:

  • 必须处理HTTP响应:仅当收到2xx状态码时视为发送成功,非2xx需触发重试。
  • 建议添加请求签名:在请求头中加入自定义签名(如基于请求体和密钥生成的HMAC),让客户侧验证请求合法性,防止伪造。
C# Azure函数(定时器触发器)可运行示例

前置准备

  1. 创建Azure函数项目(选择.NET 6+/Isolated Worker模型)。
  2. 安装NuGet包:Microsoft.Azure.Functions.Worker、Microsoft.Data.SqlClient、System.Text.Json。
  3. 在Azure函数应用的应用设置中添加:
    • DbConnectionString:你的数据库连接字符串。
    • CustomerWebhookUrl:客户提供的Webhook接收URL。

示例代码

using Microsoft.Azure.Functions.Worker;
using Microsoft.Extensions.Logging;
using Microsoft.Data.SqlClient;
using System.Text.Json;
using System.Net.Http;
using System.Threading.Tasks;

namespace DbChangeWebhookFunction
{
    public class LastNameChangeNotifier
    {
        private readonly HttpClient _httpClient;
        private readonly ILogger<LastNameChangeNotifier> _logger;
        private readonly string _dbConnString = Environment.GetEnvironmentVariable("DbConnectionString")!;
        private readonly string _webhookUrl = Environment.GetEnvironmentVariable("CustomerWebhookUrl")!;

        // 持久化最后轮询时间(示例用静态变量,生产需替换为Azure存储/数据库表)
        private static DateTime _lastPollTimestamp = DateTime.UtcNow.AddMinutes(-5);

        public LastNameChangeNotifier(HttpClient httpClient, ILogger<LastNameChangeNotifier> logger)
        {
            _httpClient = httpClient;
            _logger = logger;
        }

        // 定时器触发器:每5分钟执行一次(cron表达式:0 */5 * * * *)
        [Function("LastNameChangeNotifier")]
        public async Task Run([TimerTrigger("0 */5 * * * *")] TimerInfo timer)
        {
            _logger.LogInformation("Starting change check at {CurrentTime}", DateTime.Now);

            // 1. 查询最近轮询后的姓氏变更记录
            var changeRecords = await FetchRecentLastNameChanges();

            if (changeRecords.Count == 0)
            {
                _logger.LogInformation("No pending last name changes found.");
                return;
            }

            // 2. 批量发送Webhook通知
            foreach (var record in changeRecords)
            {
                await SendChangeNotification(record);
            }

            // 3. 更新最后轮询时间(生产环境需持久化到存储)
            _lastPollTimestamp = DateTime.UtcNow;
        }

        private async Task<List<ChangeRecord>> FetchRecentLastNameChanges()
        {
            var records = new List<ChangeRecord>();

            using var connection = new SqlConnection(_dbConnString);
            await connection.OpenAsync();

            // 假设业务表包含Id、FirstName、LastName、OldLastName、LastModified字段
            var query = @"
                SELECT Id, FirstName, LastName AS NewLastName, OldLastName, LastModified
                FROM CustomerTable
                WHERE LastModified > @LastPollTime
                  AND OldLastName != LastName";

            using var command = new SqlCommand(query, connection);
            command.Parameters.AddWithValue("@LastPollTime", _lastPollTimestamp);

            using var reader = await command.ExecuteReaderAsync();
            while (await reader.ReadAsync())
            {
                records.Add(new ChangeRecord
                {
                    Id = reader.GetInt32(0),
                    FirstName = reader.GetString(1),
                    NewLastName = reader.GetString(2),
                    OldLastName = reader.GetString(3),
                    ModifiedTime = reader.GetDateTime(4)
                });
            }

            return records;
        }

        private async Task SendChangeNotification(ChangeRecord record)
        {
            try
            {
                // 构造Webhook请求体
                var webhookPayload = JsonSerializer.Serialize(new
                {
                    Event = "LastNameUpdated",
                    OccurredAt = record.ModifiedTime,
                    Payload = new
                    {
                        CustomerId = record.Id,
                        FirstName = record.FirstName,
                        PreviousLastName = record.OldLastName,
                        CurrentLastName = record.NewLastName
                    }
                });

                var requestContent = new StringContent(webhookPayload, System.Text.Encoding.UTF8, "application/json");
                // 可选:添加签名头验证请求合法性
                // requestContent.Headers.Add("X-Webhook-Signature", GenerateHmacSignature(webhookPayload));

                var response = await _httpClient.PostAsync(_webhookUrl, requestContent);
                response.EnsureSuccessStatusCode(); // 非2xx状态码抛出异常

                _logger.LogInformation("Successfully notified change for customer {Id}", record.Id);
            }
            catch (HttpRequestException ex)
            {
                _logger.LogError(ex, "Failed to send notification for customer {Id}", record.Id);
                // 生产环境建议用Polly实现指数退避重试
            }
        }

        // 变更记录模型
        private class ChangeRecord
        {
            public int Id { get; set; }
            public string FirstName { get; set; } = string.Empty;
            public string OldLastName { get; set; } = string.Empty;
            public string NewLastName { get; set; } = string.Empty;
            public DateTime ModifiedTime { get; set; }
        }

        // 可选:生成HMAC签名
        // private string GenerateHmacSignature(string payload)
        // {
        //     var secret = Environment.GetEnvironmentVariable("WebhookSecret")!;
        //     using var hmac = new HMACSHA256(Encoding.UTF8.GetBytes(secret));
        //     var hash = hmac.ComputeHash(Encoding.UTF8.GetBytes(payload));
        //     return Convert.ToBase64String(hash);
        // }
    }
}

生产环境优化建议

  1. 持久化轮询时间:将_lastPollTimestamp存储到Azure Table Storage、SQL表或Redis中,避免函数重启后重复发送通知。
  2. 重试机制:引入Polly库实现指数退避重试,处理网络波动或客户服务不可用的情况。
  3. 错误日志与监控:集成Azure Application Insights,监控Webhook发送成功率,及时排查故障。
  4. 实时性升级:若需秒级通知,切换到CDC方案:
    • 开启Azure SQL的变更数据捕获功能。
    • 配置Azure Event Grid订阅SQL变更事件,触发Azure函数。
    • 函数中解析CDC日志,直接发送Webhook。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 04:31:01