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

Kotlin callbackFlow中awaitClose未按预期触发问题排查

问题:callbackFlow的awaitClose未触发,无法注销监听器

我尝试用callbackFlow包装回调并转换为Flow,预期Flow完成时触发awaitClose代码块注销监听器,但awaitClose未按预期执行。以下是代码及运行日志:

代码

fun main() {
   runBlocking {
       println("runBlocking on ${Thread.currentThread()}")
       getLocationFlow()
           .collect { location ->
               println("Collect location [${location.lat}, ${location.lon}] on ${Thread.currentThread()}")
           }
       println("runBlocking end")
   }
}

fun getLocationFlow(): Flow<Location> {
    return callbackFlow<Location> {
        val locationListener = object : LocationListener {
            override fun onLocationUpdate(location: Location) {
                println("onLocationUpdate: [${location.lat}, ${location.lon}] on ${Thread.currentThread()}")
                trySend(location)
            }

        }
        LocationManager.registerForLocation(locationListener)

        // 未按预期执行
        awaitClose {
            LocationManager.unregisterForLocation(locationListener)
            println("awaitClose called")
        }
    }.flowOn(Dispatchers.Default)
}

object LocationManager {

    fun registerForLocation(locationListener: LocationListener) {
        println("registerForLocation on ${Thread.currentThread()}")
        (0..5).forEach {
            locationListener.onLocationUpdate(Location(it.toLong(), it*2L))
        }
    }

    fun unregisterForLocation(locationListener: LocationListener) {
        println("unregisterForLocation on ${Thread.currentThread()}")
    }
}

data class Location(
    val lat: Long,
    val lon: Long
)

interface LocationListener {
    fun onLocationUpdate(location: Location)
}

运行日志

runBlocking on Thread[main,5,main] 
registerForLocation on Thread[DefaultDispatcher-worker-1,5,main] 
onLocationUpdate: [0, 0] on Thread[DefaultDispatcher-worker-1,5,main] 
onLocationUpdate: [1, 2] on Thread[DefaultDispatcher-worker-1,5,main] 
onLocationUpdate: [2, 4] on Thread[DefaultDispatcher-worker-1,5,main] 
onLocationUpdate: [3, 6] on Thread[DefaultDispatcher-worker-1,5,main] 
onLocationUpdate: [4, 8] on Thread[DefaultDispatcher-worker-1,5,main] 
onLocationUpdate: [5, 10] on Thread[DefaultDispatcher-worker-1,5,main] 
Collect location [0, 0] on Thread[main,5,main] 
Collect location [1, 2] on Thread[main,5,main] 
Collect location [2, 4] on Thread[main,5,main] 
Collect location [3, 6] on Thread[main,5,main] 
Collect location [4, 8] on Thread[main,5,main] 
Collect location [5, 10] on Thread[main,5,main]

问题原因

  1. Flow未主动关闭:callbackFlow的awaitClose只会在Flow通道关闭或取消时执行。你的代码中,LocationManager.registerForLocation是同步阻塞执行,一次性发送完所有数据后,没有调用close()关闭Flow通道,导致collect会一直等待新数据,永远不会结束,awaitClose也就没机会执行。
  2. 同步回调的局限性:示例中的registerForLocation是同步循环调用回调,不符合真实场景中异步持续回调的逻辑,也让Flow没有机会进入“等待关闭”的状态。

解决方法

场景1:一次性获取定位(如示例)

在发送完所有数据后,手动调用close()关闭Flow通道,这样collect会正常结束,进而触发awaitClose:

fun getLocationFlow(): Flow<Location> {
    return callbackFlow<Location> {
        val locationListener = object : LocationListener {
            override fun onLocationUpdate(location: Location) {
                println("onLocationUpdate: [${location.lat}, ${location.lon}] on ${Thread.currentThread()}")
                trySend(location)
            }
        }
        LocationManager.registerForLocation(locationListener)
        
        // 关闭通道,告知Flow无更多数据
        close()

        awaitClose {
            LocationManager.unregisterForLocation(locationListener)
            println("awaitClose called")
        }
    }.flowOn(Dispatchers.Default)
}

场景2:持续异步定位(真实场景)

修改LocationManager模拟异步持续回调,然后通过取消订阅触发awaitClose:

第一步:修改LocationManager为异步调用

object LocationManager {
    fun registerForLocation(locationListener: LocationListener) {
        println("registerForLocation on ${Thread.currentThread()}")
        // 用协程模拟异步持续发送定位
        CoroutineScope(Dispatchers.Default).launch {
            (0..5).forEach {
                delay(100) // 模拟定位间隔
                locationListener.onLocationUpdate(Location(it.toLong(), it*2L))
            }
        }
    }

    fun unregisterForLocation(locationListener: LocationListener) {
        println("unregisterForLocation on ${Thread.currentThread()}")
    }
}

第二步:主动取消订阅触发awaitClose

fun main() {
   runBlocking {
       println("runBlocking on ${Thread.currentThread()}")
       // 用Job控制Flow的生命周期
       val collectJob = launch {
           getLocationFlow().collect { location ->
               println("Collect location [${location.lat}, ${location.lon}] on ${Thread.currentThread()}")
           }
       }
       delay(700) // 等待所有定位数据发送完成
       collectJob.cancel() // 取消订阅,触发awaitClose
       println("runBlocking end")
   }
}

修改后的运行日志(场景1)

runBlocking on Thread[main,5,main]
registerForLocation on Thread[DefaultDispatcher-worker-1,5,main]
onLocationUpdate: [0, 0] on Thread[DefaultDispatcher-worker-1,5,main]
onLocationUpdate: [1, 2] on Thread[DefaultDispatcher-worker-1,5,main]
onLocationUpdate: [2, 4] on Thread[DefaultDispatcher-worker-1,5,main]
onLocationUpdate: [3, 6] on Thread[DefaultDispatcher-worker-1,5,main]
onLocationUpdate: [4, 8] on Thread[DefaultDispatcher-worker-1,5,main]
onLocationUpdate: [5, 10] on Thread[DefaultDispatcher-worker-1,5,main]
Collect location [0, 0] on Thread[main,5,main]
Collect location [1, 2] on Thread[main,5,main]
Collect location [2, 4] on Thread[main,5,main]
Collect location [3, 6] on Thread[main,5,main]
Collect location [4, 8] on Thread[main,5,main]
Collect location [5, 10] on Thread[main,5,main]
awaitClose called
unregisterForLocation on Thread[DefaultDispatcher-worker-1,5,main]
runBlocking end

内容的提问来源于stack exchange,提问作者neo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 14:27:11