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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 00:32:03