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
相关产品推荐
相关产品推荐

