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

Kotlin协程中热流的发射与订阅协调问题

问题描述

我正尝试设计一个具备以下特性的可观察任务类:

  • 响应式报告当前状态变化
  • 共享状态与结果事件:新订阅者能收到订阅后的变更通知
  • 拥有基于CoroutineScope的生命周期
  • 接口中无挂起函数(因自带生命周期)

基础代码如下:

class Worker {
  enum class State { Running, Idle }
  private val state = MutableStateFlow(State.Idle)
  private val results = MutableSharedFlow<String>()
  private val scope = CoroutineScope(Dispatchers.Default)

  private suspend fun doWork(): String {
    println("doing work")
    return "Result of the work"
  }

  fun start() {
    scope.launch {
      state.value = State.Running

      results.emit(doWork())

      state.value = State.Idle
    }
  }

  fun state(): Flow<State> = state

  fun results(): Flow<String> = results
}

当尝试“在订阅后启动工作”时出现问题,目前没有清晰的实现方式。如下简单写法无法正常运行:

fun main() {
  runBlocking {
    val worker = Worker()
    // subscriber 1
    launch {
      worker.results().collect { println("received result $it") }
    }
    worker.start()
    // subscriber 2 can also be created "later" and watch
    // for state()/result() changes
  }
}

该代码仅打印“doing work”,却无法打印结果。我理解原因在于collect和start处于不同协程中,未做任何同步处理。在doWork的协程中添加delay(300)可“解决”问题,能正常打印结果,但我希望无需人工延迟即可实现。另一种“方案”是基于results()创建SharedFlow并使用其onSubscription调用start(),但上次尝试后也未成功。

我的问题:

  1. 该设计能否修复以正常工作,还是从一开始就存在缺陷?
  2. 若设计存在缺陷,是否有其他仍能满足上述所有特性的实现方案?

解决方案

问题1:当前设计可修复,核心问题在于协程生命周期与SharedFlow配置

当前设计并非从根上存在缺陷,问题出在两个关键点:

  1. Worker内部的CoroutineScope独立于调用者的作用域(比如runBlocking),导致runBlocking不会等待Worker的协程完成,程序提前退出时协程被取消,事件无法发送。
  2. MutableSharedFlow默认配置(replay=0、extraBufferCapacity=0)下,若emit执行时无订阅者准备就绪,事件会丢失或协程挂起,最终因程序退出被取消。

修复方案1:调整SharedFlow配置并同步生命周期

修改results的SharedFlow配置,设置replay=1确保新订阅者能收到最近的结果;同时添加方法让外部等待Worker的工作完成:

class Worker {
  enum class State { Running, Idle }
  private val state = MutableStateFlow(State.Idle)
  // 设置replay=1,确保新订阅者能获取最近的结果
  private val results = MutableSharedFlow<String>(replay = 1)
  private val scope = CoroutineScope(Dispatchers.Default)

  private suspend fun doWork(): String {
    println("doing work")
    return "Result of the work"
  }

  fun start() {
    scope.launch {
      state.value = State.Running
      results.emit(doWork())
      state.value = State.Idle
    }
  }

  fun state(): Flow<State> = state
  fun results(): Flow<String> = results
  
  // 让外部等待所有工作完成,无需暴露内部协程细节
  fun awaitCompletion() = runBlocking {
    scope.coroutineContext.job.join()
  }
}

对应的main函数调用:

fun main() {
  runBlocking {
    val worker = Worker()
    // 启动订阅者
    launch {
      worker.results().collect { println("received result $it") }
    }
    worker.start()
    // 等待Worker工作完成,避免程序提前退出
    worker.awaitCompletion()
    // 晚启动的订阅者也能获取到结果
    launch {
      worker.results().collect { println("subscriber 2 received $it") }
    }
  }
}

修复方案2:等待订阅者就绪后再执行任务

如果不需要保留历史结果,可以在start中等待至少一个订阅者订阅后再执行工作,利用SharedFlow的subscriptionCount:

fun start() {
  scope.launch {
    // 等待至少一个订阅者就绪
    results.subscriptionCount.first { it > 0 }
    state.value = State.Running
    results.emit(doWork())
    state.value = State.Idle
  }
}

这种方式确保emit执行时已有订阅者,事件不会丢失,且无需修改SharedFlow的replay配置。


问题2:更健壮的替代实现方案

以下实现完全满足所有需求,同时提升了健壮性:

class Worker(private val parentScope: CoroutineScope = CoroutineScope(Dispatchers.Default)) {
  enum class State { Running, Idle }
  // 用asStateFlow暴露只读状态流
  private val _state = MutableStateFlow(State.Idle)
  val state: Flow<State> = _state.asStateFlow()
  
  // 配置replay=1,确保新订阅者能获取最近的结果
  private val _results = MutableSharedFlow<String>(replay = 1)
  val results: Flow<String> = _results.asSharedFlow()

  private suspend fun doWork(): String {
    println("doing work")
    return "Result of the work"
  }

  fun start() {
    parentScope.launch {
      // 避免重复启动任务
      if (_state.value == State.Running) return@launch
      _state.value = State.Running
      try {
        _results.emit(doWork())
      } finally {
        // 确保任务结束后状态恢复为Idle,即使任务被取消
        _state.value = State.Idle
      }
    }
  }

  // 提供取消方法,控制Worker生命周期
  fun cancel() {
    parentScope.cancel()
  }
}

特性说明

  1. 响应式状态报告:使用StateFlow管理状态,订阅者能实时获取状态变化。
  2. 共享事件与状态:results流设置replay=1,新订阅者能收到最近的结果;state流默认会发送当前状态给新订阅者。
  3. CoroutineScope生命周期:支持传入父作用域,也可使用默认作用域,生命周期由作用域控制,调用cancel()可终止所有任务。
  4. 无挂起函数接口:所有协程逻辑内部处理,对外暴露的接口均为非挂起函数。

使用示例

fun main() = runBlocking {
  // 传入runBlocking的作用域,确保程序等待Worker任务完成
  val worker = Worker(this)
  // 第一个订阅者
  launch {
    worker.results.collect { println("received result $it") }
  }
  worker.start()
  // 延迟启动第二个订阅者
  delay(100)
  launch {
    worker.results.collect { println("subscriber 2 received $it") }
  }
  // 监听状态变化,状态变为Idle后取消收集
  launch {
    worker.state.collect {
      println("current state: $it")
      if (it == Worker.State.Idle) this@launch.cancel()
    }
  }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 07:50:28