Google BigQuery的InsertRowsAsync()无法插入数据问题求解
问题诱因
- 核心原因是异步调用存在**火后不管(fire and forget)**的逻辑错误,整个异步调用链没有正确等待任务完成:
- 外层
Run方法调用异步方法SaveQueriesToMasterTableAsync时未加await,方法不会等待插入操作完成就继续执行下一个客户的逻辑,定时任务线程结束后,未完成的插入任务会被系统回收终止,小数据量时刚好能在回收前跑完所以偶尔成功,大数据量还没执行完就被终止,因此大概率失败。 - 最初版本的
SaveQueriesToMasterTableAsync内部调用InsertRowsAsync也未加await,即使你后续给该方法加了await,外层调用未等待的问题依然存在,因此问题没有解决。
- 外层
- 次要原因是insertId重复:原来的insertId生成规则是
row+i,不同客户的批次中i值会重复,BigQuery会将insertId相同的请求判定为重复插入,直接丢弃请求避免重复数据。 - 无异常抛出是因为未等待的异步任务抛出的异常会被默认的TaskScheduler捕获,不会抛到主线程,属于静默失败,所以你看不到报错。
修复方案
- 补全整个异步调用链的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); } } }
- 完善
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; } } }
- 如果定时任务框架要求入口方法必须是void类型,新增同步入口适配:
// 定时器回调指向该方法即可 private void RunEntry() { Run().GetAwaiter().GetResult(); }
替代方案
- 数据量稳定在千级以上的场景,优先采用批量导入方案:将查询历史先导出为CSV/Parquet文件上传到对象存储,再调用BigQuery的LoadJob接口批量导入,稳定性更高,成本也比流式插入更低。
- 不想调整异步逻辑的话,直接调用
InsertRows同步方法执行插入,等待操作完成后再继续执行,逻辑更简单,仅性能略低于异步方案。 - 单次插入数据量过大时可以拆分批次,每500~1000条为一批插入,避免单次请求超时。
内容的提问来源于stack exchange,提问作者RakshithaRam
相关产品推荐
相关产品推荐

