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

Kotlin KMM中无静态变量时正确停止带无限循环的协程Flow

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.02 03:01:17