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

Google BigQuery的InsertRowsAsync()无法插入数据问题求解

问题诱因
  • 核心原因是异步调用存在**火后不管(fire and forget)**的逻辑错误,整个异步调用链没有正确等待任务完成:
    1. 外层Run方法调用异步方法SaveQueriesToMasterTableAsync时未加await,方法不会等待插入操作完成就继续执行下一个客户的逻辑,定时任务线程结束后,未完成的插入任务会被系统回收终止,小数据量时刚好能在回收前跑完所以偶尔成功,大数据量还没执行完就被终止,因此大概率失败。
    2. 最初版本的SaveQueriesToMasterTableAsync内部调用InsertRowsAsync也未加await,即使你后续给该方法加了await,外层调用未等待的问题依然存在,因此问题没有解决。
  • 次要原因是insertId重复:原来的insertId生成规则是row+i,不同客户的批次中i值会重复,BigQuery会将insertId相同的请求判定为重复插入,直接丢弃请求避免重复数据。
  • 无异常抛出是因为未等待的异步任务抛出的异常会被默认的TaskScheduler捕获,不会抛到主线程,属于静默失败,所以你看不到报错。
修复方案
  1. 补全整个异步调用链的await逻辑,将Run方法改为异步Task返回,调用SaveQueriesToMasterTableAsync时添加await关键字:
private async Task Run()
{ 
   foreach(var customer in customers)
   {
      List<QueryEntry> queryEntries = bigQueryDataAccess.GetNewQueries(customer);
      if (queryEntries.Any())
      {                    
          await bigQueryDataAccess.SaveQueriesToMasterTableAsync(queryEntries);
      }
   }
}
  1. 完善SaveQueriesToMasterTableAsync的异常捕获和insertId生成逻辑:
public class BigQueryDataAccess{ 
    public async Task SaveQueriesToMasterTableAsync(List<QueryEntry> queryEntries) 
    {
        try
        {
            BigQueryInsertRow[] rows = new BigQueryInsertRow[queryEntries.Count];
            for (int i = 0; i < queryEntries.Count; i++)
            {         
                // 不需要去重可直接省略insertId参数,需要的话用Guid保证全局唯一
                BigQueryInsertRow row = new BigQueryInsertRow(insertId: $"{Guid.NewGuid()}_{i}"){
                    {"customer_id", queryEntries[i].CustomerId },
                    {"query", queryEntries[i].Query },
                    {"start_time", queryEntries[i].StartTime},
                    {"end_time", queryEntries[i].EndTime}                 
                };
                rows[i] = row;        
            }     
            await _bigQueryClientMaster.InsertRowsAsync("dataset_id", "table_id", rows);
        }
        catch(Exception ex)
        {
            // 自行替换为项目的日志组件输出异常
            Console.WriteLine($"BigQuery插入失败: {ex}", ex);
            throw;
        }
    }
}
  1. 如果定时任务框架要求入口方法必须是void类型,新增同步入口适配:
// 定时器回调指向该方法即可
private void RunEntry()
{
    Run().GetAwaiter().GetResult();
}
替代方案
  • 数据量稳定在千级以上的场景,优先采用批量导入方案:将查询历史先导出为CSV/Parquet文件上传到对象存储,再调用BigQuery的LoadJob接口批量导入,稳定性更高,成本也比流式插入更低。
  • 不想调整异步逻辑的话,直接调用InsertRows同步方法执行插入,等待操作完成后再继续执行,逻辑更简单,仅性能略低于异步方案。
  • 单次插入数据量过大时可以拆分批次,每500~1000条为一批插入,避免单次请求超时。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 06:54:05