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]
问题原因
- Flow未主动关闭:
callbackFlow的awaitClose只会在Flow通道关闭或取消时执行。你的代码中,LocationManager.registerForLocation是同步阻塞执行,一次性发送完所有数据后,没有调用close()关闭Flow通道,导致collect会一直等待新数据,永远不会结束,awaitClose也就没机会执行。 - 同步回调的局限性:示例中的
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
相关产品推荐
相关产品推荐

