如何在runBlocking中取消协程并让程序正常退出?
问题描述
开发了一段从数据库扫描行数据并通过Producer Channel批量输出的代码,期望并发处理批量数据、监控错误,在满足条件或出错时取消生产者通道与数据处理器。但调用cancel()后程序挂起,无法正常退出runBlocking,也无法输出After runBlocking.。尝试过保存runBlocking的coroutineContext并取消、取消子协程、添加yield()等操作,均未解决问题。
示例代码:
fun main(args: Array<String>) { println("Starting runBlocking.") try { runBlocking { // produceData is my producer channel val data = produceData() var count = 0 data.consumeEach { batch -> yield() count++ if (count >= 10) { println("Reached 10th batch, cancelling...") // cancel here cancel() } } } } catch (e: Throwable) { println("Exception encountered.") println(e.stackTraceToString()) } println("After runBlocking.") } @OptIn(ExperimentalCoroutinesApi::class) fun CoroutineScope.produceData(): ReceiveChannel<List<Data>> { return produce { db.connection.use { conn -> conn.createStatement(ResultSet.TYPE_FORWARD_ONLY, ResultSet.CONCUR_READ_ONLY).use { stmt -> stmt.fetchSize = Int.MIN_VALUE val sql = // omitted working SQL stmt.executeQuery(sql).use { rs -> val batch = mutableListOf<Data>() while (rs.next()) { // Handle result set here val data = Data(...) // omitted model mapper batch.add(data) if (batch.size >= 100) { send(batch) batch.clear() } } } } } } }
当前程序输出:
Starting runBlocking. Reached 10th batch, cancelling...
之后程序挂起,无法继续执行。
问题原因
核心问题在于生产者协程中的JDBC阻塞操作无法响应协程取消:
- 调用
cancel()后,协程的取消状态被标记,但JDBC的rs.next()是阻塞式IO调用,不会主动检查协程的isActive状态,导致生产者协程一直卡在数据库结果集的遍历循环中,无法正常退出。 runBlocking会等待所有子协程(包括生产者协程)完成后才会结束,因此生产者协程挂起会导致整个runBlocking无法退出。
解决方案
需要在生产者协程的数据库遍历循环中主动检查协程取消状态,一旦检测到取消就终止循环并释放资源,具体修改如下:
修改produceData函数
在rs.next()的循环内,每次迭代先检查协程是否还活跃,若已取消则立即跳出循环:
@OptIn(ExperimentalCoroutinesApi::class) fun CoroutineScope.produceData(): ReceiveChannel<List<Data>> { return produce { db.connection.use { conn -> conn.createStatement(ResultSet.TYPE_FORWARD_ONLY, ResultSet.CONCUR_READ_ONLY).use { stmt -> stmt.fetchSize = Int.MIN_VALUE val sql = // omitted working SQL stmt.executeQuery(sql).use { rs -> val batch = mutableListOf<Data>() while (rs.next()) { // 主动检查协程是否已取消,若取消则终止循环 if (!isActive) { break } // Handle result set here val data = Data(...) // omitted model mapper batch.add(data) if (batch.size >= 100) { send(batch) batch.clear() } } } } } } }
额外优化(可选)
如果数据库结果集遍历耗时极长,可以在循环中加入yield(),让协程有机会处理取消信号:
while (rs.next()) { if (!isActive) break yield() // 让出执行权,处理取消信号 // ... 数据映射逻辑 }
为什么之前的操作无效?
- 直接调用
cancel()只是标记协程取消状态,但不会中断阻塞的JDBC调用,必须主动检查isActive才能终止循环。 yield()在消费者协程中添加无效,因为问题出在生产者协程的阻塞操作上,需要在生产者的循环中添加取消检查。
内容的提问来源于stack exchange,提问作者nathlrowe
相关产品推荐
相关产品推荐

