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?>

更多推荐