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

ASP.NET Core Web API高并发下BigQuery Schema更新异步等待问题排查

现有实现存在的核心问题
  • 非线程安全容器+无互斥判断逻辑引发竞态:你使用的普通Dictionary<string, Task>本身不支持并发读写,高并发场景下同时读写、同key重复赋值会出现数据损坏、任务引用丢失的问题。同时IsUpdateNeeded判断和更新任务启动两个步骤没有互斥保护,多个并发请求可能同时判定需要更新Schema,重复触发多个更新任务,后启动的任务会覆盖字典中已有的任务引用,被覆盖的任务如果因资源不足被取消,所有等待该任务的请求就会直接抛出取消异常。
  • 自旋等待逻辑引发线程池耗尽死锁:WriteEventAsync里用while(TryGetValue + Task.Yield())实现的自旋等待,在高并发下会产生巨量无效线程调度:数百个请求同时空转轮询字典,会占满线程池工作线程,导致真正需要执行的Schema更新任务拿不到线程资源无法推进,形成“等待的请求占着线程、更新任务拿不到线程跑不完”的死锁,最终等待超时触发任务取消异常。另外如果首个写入请求还没来得及把任务存入字典,所有后续请求会进入无限自旋,永远无法继续执行。
  • 更新任务无生命周期管理,故障无法恢复:你把更新任务存入字典后从未做清理,一旦某次更新任务因异常、取消进入终态,后续所有请求都会永久拿到这个故障任务直接抛出异常,服务永远无法自动恢复。同时你没有校验字典中已存任务的状态,如果已存任务已经完成,后续请求还会无谓地等待一个已经结束的任务,甚至因为已结束任务的取消状态直接报错。
  • 更新触发与写入流程无强时序保证:你把Schema更新的判断放在SendAsync里,等待逻辑放在WriteEventAsync里,两个步骤之间没有原子性约束:可能出现请求判定不需要更新,直接进入写入流程时,首个触发更新的请求还没把任务存入字典,导致写入逻辑拿不到任务进入自旋;也可能出现更新任务还没提交到BigQuery,写入请求已经发出去,因为Schema不匹配写入失败。
正确实现思路

核心是把“Schema更新检查-更新任务启动-等待更新完成”整个逻辑做成单表粒度的原子操作,彻底消除竞态和无效等待:

  • 替换线程不安全容器:把普通Dictionary换成ConcurrentDictionary<string, Task>存储单表正在运行的Schema更新任务,所有对字典的操作都用并发字典自带的原子方法,避免读写冲突。
  • 移除自旋等待:去掉WriteEventAsync里的轮询逻辑,把Schema校验和等待的逻辑统一前置到写入操作之前,所有等待都用await实现,不阻塞工作线程,避免线程池耗尽。
  • 加单表粒度的异步互斥锁:为每个表维护一个初始计数为1的SemaphoreSlim,需要触发Schema更新时先获取锁,拿到锁后做二次校验(双检锁模式),避免重复提交更新任务。
  • 完善任务生命周期管理:更新任务无论执行成功、失败还是取消,都要在执行结束后从字典中移除,避免故障任务永久残留导致服务不可恢复;任务执行失败后,后续请求可以重新触发新的更新任务。

核心实现代码参考:

// 并发存储单表的更新任务、单表异步锁
private readonly ConcurrentDictionary<string, Task> _schemaUpdateTasks = new();
private readonly ConcurrentDictionary<string, SemaphoreSlim> _tableLocks = new();

private async Task EnsureSchemaReady(string tableName, Data sampleData)
{
    // 不需要更新Schema时,只需要等待已有运行中的更新任务完成即可
    if (!_schemaManager.IsUpdateNeeded(sampleData))
    {
        if (_schemaUpdateTasks.TryGetValue(tableName, out var runningTask))
            await runningTask;
        return;
    }

    // 获取单表专属的异步锁,保证同一张表同一时间只有一个请求在执行Schema更新检查
    var tableLock = _tableLocks.GetOrAdd(tableName, _ => new SemaphoreSlim(1, 1));
    await tableLock.WaitAsync();
    try
    {
        // 双检:拿到锁后再次判断是否需要更新,避免锁等待期间其他请求已经完成了更新
        if (!_schemaManager.IsUpdateNeeded(sampleData))
        {
            if (_schemaUpdateTasks.TryGetValue(tableName, out var runningTask))
                await runningTask;
            return;
        }

        // 原子性启动更新任务,存入字典
        var updateTask = _schemaUpdateTasks.GetOrAdd(tableName, async _ =>
        {
            try
            {
                var newSchema = _schemaManager.GetSchema(sampleData);
                await _bigQuery.UpdateTableSchemaAsync(newSchema);
            }
            finally
            {
                // 无论更新成功失败,任务结束后立刻从字典移除,避免故障残留
                _schemaUpdateTasks.TryRemove(tableName, out _);
            }
        });
        await updateTask;
    }
    finally
    {
        tableLock.Release();
    }
}

调整SendAsync逻辑,把Schema校验前置到数据准备和写入之前:

public async Task SendAsync(List<Data> data)
{
    // 先保证Schema就绪,再做后续操作
    await EnsureSchemaReady(_tableName, data.First());
    // prepare data and do some related stuff
    await _bigQuery.WriteEventAsync(preparedData);
}

调整后WriteEventAsync不需要再做任何Schema相关的等待和校验,彻底移除原有的自旋轮询逻辑即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 12:57:14