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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 09:35:22