Sql Dependency订阅user_actions表时出现行数据遗漏问题求助
问题原因分析
- 事件合并触发导致漏处理:SqlDependency不会为每一条插入操作单独触发
OnChange事件,短时间内的批量插入会被SQL Server合并为一次通知。你当前每次只查询最新的一条id,必然会漏掉中间的所有新增行。 - 订阅恢复不及时:在
OnChange事件中处理完单条数据后才重新初始化订阅,这段时间窗口内如果有新的插入,会因为没有活跃的SqlDependency订阅而错过通知。 - 数据获取逻辑错误:
select top 1 id from user_actions order by id desc只能拿到最后一条数据,完全无法覆盖批量插入的场景。
修复方案
1. 跟踪已处理的最大ID,批量获取未处理数据
维护一个全局变量记录上一次处理到的最大ID,每次触发通知时,查询所有大于该ID的新增行,而不是只取最新一条。
2. 优先恢复订阅,避免遗漏新通知
在OnChange事件的开头就重新初始化SqlDependency订阅,取消旧订阅避免内存泄漏,确保处理数据的过程中不会错过新的变更通知。
3. 优化查询与处理逻辑
确保查询符合SqlDependency的要求(已指定dbo schema,此点无需修改),同时将数据处理逻辑改为批量处理,避免单条处理的局限性。
修复后的示例代码
private long _lastProcessedId = 0; // 记录最后处理的ID,初始值设为表中最大ID或0 void Initialization() { SqlDependency.Stop(connectionString); SqlDependency.Start(connectionString); // 初始化时获取当前最大ID,避免重复处理历史数据 var maxId = Helper.ExecuteScalar("SELECT ISNULL(MAX(id), 0) FROM dbo.user_actions"); if (maxId != DBNull.Value) { _lastProcessedId = Convert.ToInt64(maxId); } SqlDependencyInit(); } void SqlDependencyInit() { using (var conn = new SqlConnection(connectionString)) { conn.Open(); using (SqlCommand command = conn.CreateCommand()) { command.CommandType = CommandType.Text; command.CommandText = "SELECT id FROM dbo.user_actions"; command.Notification = null; SqlDependency dependency = new SqlDependency(command); dependency.OnChange += OnDependencyChange; // 执行命令激活订阅 using (command.ExecuteReader()) { // 无需处理结果,仅用于激活订阅 } } } } void OnDependencyChange(object sender, SqlNotificationEventArgs e) { // 先取消旧订阅,避免内存泄漏 ((SqlDependency)sender).OnChange -= OnDependencyChange; // 立即重新初始化订阅,确保不遗漏新通知 SqlDependencyInit(); if (e.Info == SqlNotificationInfo.Insert) { // 查询所有未处理的新增行 var unprocessedActions = Helper.ExecuteDataTable($"SELECT id, action FROM dbo.user_actions WHERE id > {_lastProcessedId} ORDER BY id ASC"); foreach (DataRow row in unprocessedActions.Rows) { long id = Convert.ToInt64(row["id"]); processUserAction(id); _lastProcessedId = id; // 更新最后处理的ID } } }
额外注意事项
- 确保
Helper.ExecuteDataTable等工具方法是线程安全的,因为OnChange事件在后台线程触发。 - 如果表的
id不是自增主键,可改用date_created字段跟踪未处理数据(记录最后处理时间,查询date_created > @LastProcessedTime),但自增ID的方式可靠性更高。 - 若
processUserAction操作耗时,建议将数据放入队列,用单独线程处理,避免阻塞订阅恢复流程。
内容的提问来源于stack exchange,提问作者Jem
相关产品推荐
相关产品推荐

