如何编写可提前唤醒的Kotlin Flow/协程延迟逻辑
解决方案
要实现调用refresh()时立即中断当前延迟并重新执行认证状态检查,核心思路是用信号流触发刷新操作,在Flow循环中同时监听延迟和刷新信号,任一条件满足就进入下一次循环。
修改后的完整代码
import kotlinx.coroutines.delay import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.MutableSharedFlow import kotlinx.coroutines.flow.first import kotlinx.coroutines.flow.flow import kotlinx.coroutines.withTimeoutOrNull import java.util.Date sealed class AuthState { object UnAuthenticated : AuthState() data class Authenticated(val expirationDate: Date) : AuthState() } class AuthRepo(private val authData: AuthData) { // 刷新信号流,extraBufferCapacity=1确保未处理的最新信号不会丢失 private val refreshTrigger = MutableSharedFlow<Unit>(extraBufferCapacity = 1) fun getAuthFlow(): Flow<AuthState> = flow { while (currentCoroutineContext().isActive) { val expireDate = authData.getAuthExpirationDate() val now = Date() // 发射当前认证状态 val currentState = if (expireDate.before(now)) { AuthState.UnAuthenticated } else { AuthState.Authenticated(expireDate) } emit(currentState) // 计算延迟时长:未认证用30秒,已认证用过期剩余时间 val delayAmt = if (expireDate.before(now)) { MAX_DELAY } else { expireDate.time - now.time } // 等待延迟完成,或收到刷新信号就立即进入下一次循环 withTimeoutOrNull(delayAmt) { refreshTrigger.first() } } } fun refresh() { // 非协程环境下安全发送刷新信号 refreshTrigger.tryEmit(Unit) } companion object { private const val MAX_DELAY = 30000L // 30秒 } } // 假设的认证数据接口 interface AuthData { fun getAuthExpirationDate(): Date }
关键细节说明
MutableSharedFlow信号流:设置extraBufferCapacity=1,保证即使Flow还在处理上一轮循环,新的刷新信号也能被保留,不会丢失。tryEmit支持在非协程线程调用,适配外部触发的场景。withTimeoutOrNull+refreshTrigger.first():这个组合实现了“等待指定时长,或收到刷新信号立即返回”的逻辑。延迟到期时withTimeoutOrNull返回null;收到刷新信号则提前结束等待,直接进入下一轮状态检查。- 循环活跃性控制:
currentCoroutineContext().isActive确保Flow被取消时能及时终止循环,避免内存泄漏。
替代方案:select表达式写法
如果觉得withTimeoutOrNull不够直观,也可以用select表达式更清晰地表达“二选一”逻辑:
import kotlinx.coroutines.selects.select // 替换getAuthFlow中的等待代码块 select<Unit> { // 延迟时间到触发 onTimeout(delayAmt) {} // 收到刷新信号触发 refreshTrigger.onReceive {} }
两种方案功能完全一致,select写法更直白地展现了“等待延迟或刷新信号”的逻辑。
内容的提问来源于stack exchange,提问作者Matt Wolfe
相关产品推荐
相关产品推荐

