Kotlin KMM中无静态变量时正确停止带无限循环的协程Flow
需求背景
开发新闻获取类KMM应用,核心规则如下:
- 应用每30秒拉取一次新闻存入本地数据库,仅登录用户可使用该功能
- 用户登出时,需停止新闻刷新任务,同时清空本地数据库
- 禁止使用静态变量实现Flow的安全停止
当前项目分层架构: - ViewModel:Android、iOS端分别实现(Android用Jetpack ViewModel,iOS用ObservableObject)
- UseCase:跨端共享
- Repository:跨端共享
- 数据源:跨端共享
- Android端采用Jetpack Compose单Activity架构
- iOS端采用SwiftUI架构
现有代码实现
Android端ViewModel(iOS逻辑一致)
@HiltViewModel class NewsViewModel @Inject constructor( private val startFetchingNews: GetNewsUseCase, private val stopFetchingNews: StopGettingNewsUseCase, ) : ViewModel() { private val _mutableNewsUiState = MutableStateFlow(NewsState()) val newsUiState: StateFlow<NewsState> get() = _mutableNewsUiState.asStateFlow() fun onTriggerEvent(action: MapEvents) { when (action) { is NewsEvent.GetNews -> { getNews() } is MapEvents.StopNews -> { // 待实现停止逻辑 } else -> {} } } private fun getNews() { startFetchingNews().collectCommon(viewModelScope) { result -> when { result.error -> { // 更新UI错误状态 } result.succeeded -> { // 更新UI新闻列表状态 } } } } }
拉取新闻UseCase
class GetNewsUseCase( private val newsRepo: NewsRepoInterface) { companion object { private val UPDATE_INTERVAL = 30.seconds } operator fun invoke(): CommonFlow<Result<List<News>>> = flow { while (true) { emit(Result.loading()) val result = newsRepo.getNews() if (result.succeeded) { // emit成功结果 } else { //emit错误结果 } delay(UPDATE_INTERVAL) } }.asCommonFlow() }
新闻仓库实现
class NewsRepository( private val sourceNews: SourceNews, private val cacheNews: CacheNews) : NewsRepoInterface { override suspend fun getNews(): Result<List<News>> { val news = sourceNews.fetchNews() // 业务逻辑处理 cacheNews.insert(news) return Result.data(cacheNews.selectAll()) } }
跨端通用Flow扩展
fun <T> Flow<T>.asCommonFlow(): CommonFlow<T> = CommonFlow(this) class CommonFlow<T>(private val origin: Flow<T>) : Flow<T> by origin { fun collectCommon( coroutineScope: CoroutineScope? = null, // Android传viewModelScope,iOS传nil callback: (T) -> Unit, ) { onEach { callback(it) }.launchIn(coroutineScope ?: CoroutineScope(Dispatchers.Main)) } }
已尝试的问题方案
曾尝试将while循环下沉到Repository层,通过单例Repository中断循环,但该方案需要将getNews方法改为Flow返回类型,在UseCase层收集时会产生Flow嵌套问题,不符合分层设计要求。
最终实现方案
核心思路基于Kotlin协程结构化并发的取消机制,不需要静态变量、不需要在Flow循环内增加中断标记,仅通过协程Job的生命周期控制即可实现安全停止。
第一步:修改CommonFlow扩展,返回协程Job
原有collectCommon方法没有返回launchIn产生的Job对象,外部无法拿到协程引用执行取消,修改如下:
class CommonFlow<T>(private val origin: Flow<T>) : Flow<T> by origin { fun collectCommon( coroutineScope: CoroutineScope? = null, callback: (T) -> Unit, ): Job { // 返回Job供外部控制生命周期 return onEach { callback(it) }.launchIn(coroutineScope ?: CoroutineScope(Dispatchers.Main)) } }
第二步:修改ViewModel,持有拉取任务的Job引用
在ViewModel中增加私有变量存储拉取任务的Job,启动任务前先取消旧任务避免重复拉取,停止事件触发时直接取消Job即可终止Flow循环:
@HiltViewModel class NewsViewModel @Inject constructor( private val startFetchingNews: GetNewsUseCase, private val stopFetchingNews: StopGettingNewsUseCase, ) : ViewModel() { private val _mutableNewsUiState = MutableStateFlow(NewsState()) val newsUiState: StateFlow<NewsState> get() = _mutableNewsUiState.asStateFlow() // 持有新闻拉取协程的引用 private var fetchNewsJob: Job? = null fun onTriggerEvent(action: MapEvents) { when (action) { is NewsEvent.GetNews -> getNews() is MapEvents.StopNews -> { // 1. 取消拉取任务,Flow会自动响应协程取消终止无限循环 fetchNewsJob?.cancel() fetchNewsJob = null // 2. 执行数据库清空操作 viewModelScope.launch { stopFetchingNews() } // 3. 重置UI状态 _mutableNewsUiState.value = NewsState() } else -> {} } } private fun getNews() { // 启动新任务前先取消旧任务,避免重复创建收集器 fetchNewsJob?.cancel() fetchNewsJob = startFetchingNews().collectCommon(viewModelScope) { result -> when { result.error -> { // 更新UI错误状态 } result.succeeded -> { // 更新UI新闻列表状态 } } } } }
第三步:实现停止拉取用例与仓库清空逻辑
新增停止拉取的UseCase,仅负责调用仓库的清空缓存方法,保持单向依赖:
class StopGettingNewsUseCase( private val newsRepo: NewsRepoInterface ) { suspend operator fun invoke() { newsRepo.clearAllNewsCache() } }
在仓库接口与实现中增加清空缓存方法:
interface NewsRepoInterface { suspend fun getNews(): Result<List<News>> suspend fun clearAllNewsCache() } class NewsRepository( private val sourceNews: SourceNews, private val cacheNews: CacheNews ) : NewsRepoInterface { override suspend fun getNews(): Result<List<News>> { val news = sourceNews.fetchNews() cacheNews.insert(news) return Result.data(cacheNews.selectAll()) } override suspend fun clearAllNewsCache() { // 调用本地数据源的清空表方法 cacheNews.clearAll() } }
方案原理说明
Kotlin协程取消是协作式机制,现有UseCase中
while(true)循环内的delay()是协程原生挂起函数,会主动响应协程取消状态:当收集Flow的Job被调用cancel()时,如果处于delay等待阶段会直接抛出CancellationException终止循环;如果处于网络请求、数据库读写阶段,只要使用的是协程友好的框架(如Ktor、SQLDelight),也会自动响应取消,不会产生资源泄漏。
整个实现不需要修改原有拉取UseCase的循环逻辑,不需要静态变量存储控制位,完全遵循分层设计原则,Android、iOS端逻辑可以完全复用共享层代码。
内容的提问来源于stack exchange,提问作者Borg

