如何用C#/Azure函数实现数据库变更触发Webhook推送数据?
方案可行性分析
定时器轮询方案完全可行,但需根据你的实时性需求权衡:
- 优势:实现简单,无需额外数据库配置,适合对延迟容忍度较高(如分钟级)的场景。
- 劣势:存在轮询间隔内的延迟,若轮询过密会增加数据库负载。
如果需要近乎实时(秒级)的通知,更推荐变更数据捕获(CDC)+ 事件驱动的推模式:开启数据库CDC捕获变更事件,通过Azure Event Grid或Service Bus触发Azure函数,直接处理并推送Webhook,延迟更低且资源利用率更高。
Webhook核心实现逻辑
Webhook本质就是向客户指定的URL发起HTTP POST请求,将结构化的JSON数据作为请求体发送。关键要点:
- 必须处理HTTP响应:仅当收到2xx状态码时视为发送成功,非2xx需触发重试。
- 建议添加请求签名:在请求头中加入自定义签名(如基于请求体和密钥生成的HMAC),让客户侧验证请求合法性,防止伪造。
C# Azure函数(定时器触发器)可运行示例
前置准备
- 创建Azure函数项目(选择.NET 6+/Isolated Worker模型)。
- 安装NuGet包:
Microsoft.Azure.Functions.Worker、Microsoft.Data.SqlClient、System.Text.Json。 - 在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); // } } }
生产环境优化建议
- 持久化轮询时间:将
_lastPollTimestamp存储到Azure Table Storage、SQL表或Redis中,避免函数重启后重复发送通知。 - 重试机制:引入Polly库实现指数退避重试,处理网络波动或客户服务不可用的情况。
- 错误日志与监控:集成Azure Application Insights,监控Webhook发送成功率,及时排查故障。
- 实时性升级:若需秒级通知,切换到CDC方案:
- 开启Azure SQL的变更数据捕获功能。
- 配置Azure Event Grid订阅SQL变更事件,触发Azure函数。
- 函数中解析CDC日志,直接发送Webhook。
内容的提问来源于stack exchange,提问作者Paul
相关产品推荐
相关产品推荐

