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

ASP.NET Core 7中API端点重复调用的节流处理方案

解决ASP.NET Core 7 API节流问题:避免数据库死锁与冗余SignalR通知

问题背景

数据库发生变更时,外部服务会调用我的ASP.NET Core 7 API端点,该端点从数据库拉取更新数据并通过SignalR推送给已连接客户端。但该服务可能在几秒内发送多次调用,引发数据库死锁,还会向客户端发送冗余通知。

我需要对数据库调用和客户端通知进行节流,仅响应时间窗口内的最后一次请求,可接受1-3秒的响应延迟。目前的困惑是:每个请求处于不同线程,不知道如何实现线程间通信,也不确定使用静态属性是否合理。

补充说明

  • API调用仅需查询SQL Server并发送SignalR通知,无其他负载;
  • 端点仅返回空200 OK响应,调用方仅用于触发通知;
  • 需直接丢弃前期调用,仅处理最后一次,避免冗余数据库访问和客户端通知。

需避免的错误日志示例

RefreshOrderNotifications received at Tue Oct 03 2023 09:34:19 GMT-0400 (Eastern Daylight Time)
RefreshOrderNotifications received at Tue Oct 03 2023 09:34:57 GMT-0400 (Eastern Daylight Time)
RefreshOrderNotifications received at Tue Oct 03 2023 09:34:57 GMT-0400 (Eastern Daylight Time)
RefreshOrderNotifications received at Tue Oct 03 2023 09:34:57 GMT-0400 (Eastern Daylight Time)
Server - Error: Notification Server connection to DB failed. Transaction (Process ID 68) was deadlocked on lock resources with another process and has been chosen as the deadlock victim. Rerun the transaction.

当前实现的问题

基于参考回复我实现了ThrottleManager类,但因操作非原子性,短时间连续调用会绕过限制,代码如下:

public class ThrottleManager
{
    private List<string> _queue = new List<string>();

    public async Task<ThrottleResult> QueueThenRun(string runKey, int delayMilliseconds, Func<Task<ThrottleResult>> funcToRun)
    {
        var result = new ThrottleResult();

        if (_queue.Contains(runKey) || funcToRun == null)
        {           
            //we're just queuing up the call, so return success
            return result;
        }
        _queue.Add(runKey);

        await Task.Delay(delayMilliseconds);
        _queue.Remove(runKey);

        var returnMe = await funcToRun();
        return returnMe;
    }
}

public class ThrottleResult
{
    public bool Success { get; set; } = true;
    public string Message { get; set; } = string.Empty;
    public Exception? Exception { get; set; } = null;
}

完善方案

要实现“仅处理时间窗口内最后一次请求”的节流,需要解决线程安全问题,并支持取消之前的待执行任务。可以使用ConcurrentDictionary来跟踪每个runKey对应的待执行任务和取消令牌,确保操作原子性:

using System.Collections.Concurrent;

public class ThrottleManager
{
    // 存储每个runKey对应的取消令牌源和待完成任务
    private readonly ConcurrentDictionary<string, (CancellationTokenSource Cts, Task<ThrottleResult> Task)> _pendingTasks = new();

    public async Task<ThrottleResult> QueueThenRun(string runKey, int delayMilliseconds, Func<Task<ThrottleResult>> funcToRun)
    {
        if (funcToRun == null)
            return new ThrottleResult { Success = false, Message = "执行函数不能为空" };

        // 原子操作:获取现有任务或创建新的
        while (true)
        {
            if (_pendingTasks.TryGetValue(runKey, out var existingEntry))
            {
                // 取消之前的待执行任务,因为新请求来了
                existingEntry.Cts.Cancel();
                // 等待之前的任务完成(无论是否取消),避免内存泄漏
                try { await existingEntry.Task; } catch { }
                // 移除旧条目,准备添加新的
                _pendingTasks.TryRemove(runKey, out _);
            }
            else
            {
                var cts = new CancellationTokenSource();
                var task = ExecuteAfterDelay(runKey, delayMilliseconds, funcToRun, cts.Token);
                
                // 尝试添加新条目,如果添加失败说明有其他线程已经添加,重新循环
                if (_pendingTasks.TryAdd(runKey, (cts, task)))
                    break;
                
                // 添加失败,取消当前创建的令牌和任务
                cts.Cancel();
                try { await task; } catch { }
            }
        }

        // 直接返回成功,因为调用方只需要触发通知,不需要等待执行结果
        return new ThrottleResult();
    }

    private async Task<ThrottleResult> ExecuteAfterDelay(string runKey, int delayMilliseconds, Func<Task<ThrottleResult>> funcToRun, CancellationToken cancellationToken)
    {
        try
        {
            await Task.Delay(delayMilliseconds, cancellationToken);
            
            // 如果令牌已取消,直接返回
            if (cancellationToken.IsCancellationRequested)
                return new ThrottleResult { Success = false, Message = "任务已被取消" };

            // 执行实际操作
            var result = await funcToRun();
            
            return result;
        }
        catch (OperationCanceledException)
        {
            return new ThrottleResult { Success = false, Message = "任务已被取消" };
        }
        catch (Exception ex)
        {
            return new ThrottleResult { Success = false, Exception = ex, Message = "执行失败" };
        }
        finally
        {
            // 无论成功或失败,移除pending条目
            _pendingTasks.TryRemove(runKey, out _);
        }
    }
}

public class ThrottleResult
{
    public bool Success { get; set; } = true;
    public string Message { get; set; } = string.Empty;
    public Exception? Exception { get; set; } = null;
}

方案说明

  1. 线程安全:使用ConcurrentDictionary确保多线程环境下的原子操作,避免竞态条件;
  2. 取消旧任务:每次新请求到来时,取消之前的待执行任务,确保只有最后一次请求会触发实际操作;
  3. 内存管理:在任务完成或取消后,及时移除字典中的条目,并等待旧任务完成,避免内存泄漏;
  4. 符合需求:API端点直接返回200 OK,实际操作在延迟后执行,仅处理最后一次请求,彻底避免冗余数据库访问和SignalR通知。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 07:25:59