如何为BigQuery的InsertAll请求设置Session ID实现事务?
解决方案
BigQuery的InsertAllRequest(流式插入接口)不支持绑定Session ID,因为它属于非事务性的流式写入,无法纳入Session管理的事务中。要实现事务内的多行插入,需要改用DML INSERT语句通过QueryJobConfiguration执行,这样就能复用你已有的Session配置,将插入操作纳入事务。
代码修改步骤
- 修改
runRoutine中的调用,将事务的Session配置传递给插入方法:
transaction { configureSession -> query( """ UPDATE `${googleTableId1.toSqlTableName()}` SET $latestRefreshAt = CURRENT_TIMESTAMP() WHERE $context = @$context AND $id = @$id """, mapOf( context to QueryParameterValue.string(source), id to QueryParameterValue.string(name), ), configureSession ) // 传递configureSession参数 tableInsertRows( tableId = googleTableId2, rowContents = records, configureSession ) }
- 重写
tableInsertRows函数,用DML INSERT替代InsertAll:
suspend fun tableInsertRows( tableId: TableId, rowContents: Iterable<InsertAllRequest.RowToInsert>, configure: QueryJobConfiguration.Builder.() -> Unit = {}, ) { val rows = rowContents.map { it.content } if (rows.isEmpty()) return // 构建目标表的完整SQL名称 val fullTableName = "`${tableId.project}.${tableId.dataset}.${tableId.table}`" // 提取所有字段名(假设所有行的字段一致) val fieldNames = rows.first().keys.joinToString(", ") // 构建每行的参数占位符,避免SQL注入 val valueClauses = rows.mapIndexed { idx, _ -> fieldNames.split(", ") .map { field -> "@${field.trim()}_$idx" } .joinToString("(", ")") }.joinToString(", ") // 组装INSERT SQL val insertSql = """INSERT INTO $fullTableName ($fieldNames) VALUES $valueClauses""" // 构建命名参数映射 val parameters = rows.flatMapIndexed { idx, row -> row.map { (field, value) -> "${field}_$idx" to QueryParameterValue.of(value) } }.toMap() // 调用已有的query方法,自动绑定Session query(insertSql, parameters, configure) }
关键说明
- 流式插入(InsertAll)设计为低延迟、高吞吐量的写入,不支持事务原子性,因此无法绑定Session。
- DML INSERT语句通过QueryJob执行,属于Session事务的一部分,能保证和UPDATE操作要么同时成功,要么同时回滚。
- 如果插入的记录数量极大(比如超过1000行),建议分批次执行INSERT,避免单条SQL过长导致性能问题。
内容的提问来源于stack exchange,提问作者Gamer2015
相关产品推荐
相关产品推荐

