JDBI Kotlin PostgreSQL:如何执行@SqlBatch所有查询并获取错误行?
批量Upsert异常处理:跳过错误条目并记录异常
核心问题本质
默认JDBC批量操作(包括@SqlBatch实现的逻辑)是在同一个事务中执行所有语句,一旦某条失败会触发事务回滚,后续语句也终止执行。要实现「失败条目单独记录、其余正常提交」,需从事务隔离或批量容错机制入手。
方案1:拆分批量为单条Upsert(全数据库兼容)
将批量操作拆为单条独立执行,每条单独捕获异常,完全控制错误处理逻辑,不受数据库限制。
步骤1:扩展Repository接口
添加单条Upsert方法:
interface WriteRepository<Item> { // 单条upsert逻辑 @SqlUpdate("INSERT INTO your_table (column1, column2) VALUES (:item.column1, :item.column2) " + "ON CONFLICT (column1) DO UPDATE SET column2 = :item.column2") fun upsert(@BindBean("item") item: Item): Int // 保留原批量方法(可选) @SqlBatch("INSERT INTO your_table (column1, column2) VALUES (:item.column1, :item.column2) " + "ON CONFLICT (column1) DO UPDATE SET column2 = :item.column2") fun upsertAll(@BindBean("item") items: List<Item>): IntArray }
步骤2:业务层实现安全批量逻辑
循环处理每条数据,捕获异常并记录错误条目:
fun safeUpsertAll(items: List<Item>, repository: WriteRepository<Item>): Pair<List<Item>, List<Pair<Item, Exception>>> { val successList = mutableListOf<Item>() val errorList = mutableListOf<Pair<Item, Exception>>() items.forEach { item -> runCatching { repository.upsert(item) } .onSuccess { successList.add(item) } .onFailure { errorList.add(item to it) } } return successList to errorList }
测试调用示例
private fun testSafeBatch() { val entities = listOf( Item(1,"goodValue"), Item(2,"goodValue"), Item(3,"badValue"), Item(4,"goodValue"), Item(5,"goodValue"), Item(6,"goodValue") ) val (successes, errors) = safeUpsertAll(entities, repository) println("成功提交:${successes.size}条") println("失败条目:${errors.size}条") errors.forEach { (item, e) -> println("条目${item.column1}失败:${e.message}") } }
优缺点:
- 优势:兼容所有数据库,逻辑简单直观,完全控制错误处理
- 劣势:性能比纯批量差(每条单独提交),超大数据量需结合分批优化
方案2:利用数据库批量容错特性(部分数据库支持)
部分数据库(如PostgreSQL、MySQL)的JDBC驱动支持「批量失败时跳过错误条目,继续执行后续语句」,需配置连接参数并解析异常获取失败索引。
以PostgreSQL为例:
- 配置JDBC连接参数:
在数据源URL中添加continueBatchOnError=true:
jdbc:postgresql://localhost:5432/your_db?continueBatchOnError=true
- 修改异常处理逻辑:
通过BatchUpdateException的updateCounts数组判断失败条目(Statement.EXECUTE_FAILED即-3代表该条执行失败):
private fun testBatchWithErrorHandling() { val entities = listOf( Item(1,"goodValue"), Item(2,"goodValue"), Item(3,"badValue"), Item(4,"goodValue"), Item(5,"goodValue"), Item(6,"goodValue") ) val errorList = mutableListOf<Item>() try { val result = repository.upsertAll(entities) println("成功执行:${result.sum()}条") } catch (e: BatchUpdateException) { e.updateCounts.forEachIndexed { index, count -> if (count == Statement.EXECUTE_FAILED) { errorList.add(entities[index]) } } println("失败条目:${errorList.size}条") errorList.forEach { println("失败条目ID:${it.column1}") } // 注意:此时成功条目已提交,不会回滚 } catch (e: Exception) { println("全局异常:${e.message}") } }
注意:
- MySQL需额外配置
rewriteBatchedStatements=true,且驱动版本需支持该特性 - 该方案依赖数据库驱动,并非所有数据库都支持
方案3:混合模式(分批小批量+容错回退)
针对超大数据量场景,将数据拆分为小批量(如每50条一批),每批独立事务执行;若某批失败,则自动拆为单条处理该批次,平衡性能与容错性:
fun batchWithFallback(items: List<Item>, repository: WriteRepository<Item>, batchSize: Int = 50): Pair<List<Item>, List<Pair<Item, Exception>>> { val successList = mutableListOf<Item>() val errorList = mutableListOf<Pair<Item, Exception>>() items.chunked(batchSize).forEach { batch -> runCatching { repository.upsertAll(batch) } .onSuccess { successList.addAll(batch) } .onFailure { // 批量失败,拆为单条处理该批次 batch.forEach { item -> runCatching { repository.upsert(item) } .onSuccess { successList.add(item) } .onFailure { errorList.add(item to it) } } } } return successList to errorList }
内容的提问来源于stack exchange,提问作者Couldosh
相关产品推荐
相关产品推荐

