Kotlin协程 -> flow的lambda表达式创建Flow
·
Flow构建执行流程图
👤 用户代码调用
val myFlow = flow {
println("1. 开始构建Flow")
emit(1)
emit(2)
emit(3)
}
↓
🏗️ flow() 函数执行 - 当前线程
↓
📝 创建 SafeFlow 实例
├── 🎯 保存用户代码块: suspend FlowCollector<Int>.() -> Unit
├── 🔒 包装为 SafeFlow(block)
└── 🎁 返回 Flow<Int> 实例
↓
⏸️ Flow创建完成,等待收集 (冷流特性)
↓
👤 用户调用 collect
myFlow.collect { value ->
println("收集到: $value")
}
↓
🎯 AbstractFlow.collect() - 当前线程
↓
🛡️ 创建 SafeCollector 包装用户collector
├── 🔍 SafeCollector(userCollector, coroutineContext)
├── 🎯 检查协程上下文安全性
└── 📦 包装用户的收集逻辑
↓
🚀 调用 SafeFlow.collectSafely(safeCollector)
↓
🎯 执行用户代码块 - collector.block()
├── 📍 this = safeCollector (FlowCollector实例)
├── 🧵 println("1. 开始构建Flow")
├── 📤 emit(1) → safeCollector.emit(1)
│ ├── 🔍 SafeCollector.emit() 安全检查
│ ├── 🎯 调用用户collector: println("收集到: 1")
│ └── ✅ 第一个值发射完成
├── 📤 emit(2) → safeCollector.emit(2)
│ ├── 🔍 SafeCollector.emit() 安全检查
│ ├── 🎯 调用用户collector: println("收集到: 2")
│ └── ✅ 第二个值发射完成
└── 📤 emit(3) → safeCollector.emit(3)
├── 🔍 SafeCollector.emit() 安全检查
├── 🎯 调用用户collector: println("收集到: 3")
└── ✅ 第三个值发射完成
↓
🏁 用户代码块执行完成
↓
🧹 finally块: safeCollector.releaseIntercepted()
├── 🔧 清理拦截器资源
└── 🎯 确保协程安全结束
↓
✅ Flow收集完成
Flow构建器概览
// 1. flow构建器 - 最常用
val flow1 = flow { emit(1) }
// 2. flowOf构建器 - 固定值
val flow2 = flowOf(1, 2, 3)
// 3. asFlow扩展 - 集合转换
val flow3 = listOf(1, 2, 3).asFlow()
Flow构建器源码深度解析
1. flow函数签名和实现
/**
* 创建一个冷流,当被收集时执行给定的代码块
* @param block 在FlowCollector上下文中执行的挂起函数
* @return 新创建的Flow实例
*/
public fun <T> flow(
@BuilderInference block: suspend FlowCollector<T>.() -> Unit
): Flow<T> = SafeFlow(block)
关键点解析:
@BuilderInference:帮助编译器进行类型推断suspend FlowCollector<T>.() -> Unit:扩展函数类型,FlowCollector作为接收者- 返回
SafeFlow实例,这是Flow接口的内部实现
2. SafeFlow实现源码
/**
* SafeFlow是flow构建器的内部实现类
* 提供异常安全的Flow执行环境
*/
private class SafeFlow<T>(
private val block: suspend FlowCollector<T>.() -> Unit
) : AbstractFlow<T>() {
/**
* 安全地收集Flow数据
* @param collector 数据收集器
*/
override suspend fun collectSafely(collector: FlowCollector<T>) {
collector.block() // 在collector上下文中执行用户代码块
}
}
3. AbstractFlow基类分析
/**
* Flow接口的抽象实现,提供通用的collect方法
*/
public abstract class AbstractFlow<T> : Flow<T>, CancellableFlow<T> {
/**
* 收集Flow数据的入口方法
* @param collector 数据收集器
*/
public final override suspend fun collect(collector: FlowCollector<T>) {
val safeCollector = SafeCollector(collector, coroutineContext)
try {
collectSafely(safeCollector) // 调用子类实现
} finally {
safeCollector.releaseIntercepted() // 清理资源
}
}
/**
* 子类需要实现的安全收集方法
*/
public abstract suspend fun collectSafely(collector: FlowCollector<T>)
}
4. FlowCollector接口定义
/**
* Flow数据收集器接口
* 定义了数据发射的核心方法
*/
public fun interface FlowCollector<in T> {
/**
* 发射一个值到Flow中
* @param value 要发射的值
*/
public suspend fun emit(value: T)
}
5. emit方法调用
/**
* SafeCollector是Flow安全收集的核心实现
* 提供协程安全性检查和异常处理
*/
@PublishedApi
internal class SafeCollector<T> @JvmOverloads constructor(
@JvmField val collector: FlowCollector<T>,
@JvmField val collectContext: CoroutineContext,
@JvmField val collectContextSize: Int = collectContext.fold(0) { count, _ -> count + 1 }
) : FlowCollector<T>, ContinuationImpl(NoOpContinuation, EmptyCoroutineContext) {
// 当前协程上下文
override val context: CoroutineContext get() = completion.context
/**
* 核心emit方法实现
* 提供协程安全性检查
*/
override suspend fun emit(value: T) {
return suspendCoroutineUninterceptedOrReturn sc@{ uCont ->
try {
// 检查协程上下文安全性
emit(uCont, value)
} catch (e: Throwable) {
// 异常处理
lastEmissionContext = DownstreamExceptionElement(e)
throw e
}
}
}
/**
* 内部emit实现,包含详细的安全检查
*/
private fun emit(uCont: Continuation<Unit>, value: T): Any? {
val currentContext = uCont.context
val previousContext = lastEmissionContext
// 协程上下文切换检查
if (currentContext !== previousContext) {
checkContext(currentContext, previousContext, value)
}
completion = uCont
// 调用用户提供的collector
return emitFun(collector as FlowCollector<Any?>, value, this as Continuation<Unit>)
}
}
/**
* 实际执行emit操作的函数
* 通过函数引用优化性能
*/
private val emitFun = FlowCollector<Any?>::emit as Function3<FlowCollector<Any?>, Any?, Continuation<Unit>, Any?>
更多推荐
所有评论(0)