1. flow{} 和 callbackFlow{} 的区别

flow{} 的特点

fun simpleFlow(): Flow<String> = flow {

    // flow{} 是同步的,在协程中执行

    emit("Hello")

    delay(1000)

    emit("World")

}

callbackFlow{} 的特点

fun callbackFlowExample(): Flow<String> = callbackFlow {

    // callbackFlow{} 是异步的,可以处理回调

    // 可以在非协程环境中发送数据

}

2. 为什么 SSE 不能用 flow{}?

问题:SSE 是基于回调的

// SSE 的 HTTP 请求是基于回调的

call?.enqueue(object : Callback {

    override fun onResponse(call: Call, response: Response) {

        // 在回调中处理数据

        // 这里不能直接使用 emit()

    }

    

    override fun onFailure(call: Call, e: IOException) {

        // 错误处理

    }

})

如果强制用 flow{} 会怎样?

// 错误示例:这样不行

fun connectWithFlow(): Flow<SSEEvent> = flow {

    call?.enqueue(object : Callback {

        override fun onResponse(call: Call, response: Response) {

            // ❌ 错误:不能在回调中直接使用 emit()

            emit(SSEEvent.Message("event", "data"))

        }

    })

    // flow 会立即结束,不会等待回调

}

3. 具体对比

使用 flow{} 的问题

fun connectWithFlow(): Flow<SSEEvent> = flow {

    val request = Request.Builder()

        .url(url)

        .build()

    

    call = okHttpClient.newCall(request)

    call?.enqueue(object : Callback {

        override fun onResponse(call: Call, response: Response) {

            // ❌ 问题 1:不能在回调中使用 emit()

            // emit(SSEEvent.Message("event", "data"))

            

            // ❌ 问题 2:回调是异步的,flow 已经结束了

        }

    })

    

    // ❌ 问题 3:flow 立即结束,不会等待 HTTP 响应

}

使用 callbackFlow{} 的解决方案

fun connectWithCallbackFlow(): Flow<SSEEvent> = callbackFlow {

    val request = Request.Builder()

        .url(url)

        .build()

    

    call = okHttpClient.newCall(request)

    call?.enqueue(object : Callback {

        override fun onResponse(call: Call, response: Response) {

            // ✅ 正确:可以在回调中使用 trySend()

            trySend(SSEEvent.Message("event", "data"))

        }

        

        override fun onFailure(call: Call, e: IOException) {

            // ✅ 正确:可以发送错误

            trySend(SSEEvent.Error(e))

        }

    })

    

    // ✅ 正确:等待 Channel 关闭

    awaitClose { 

        disconnect() 

    }

}

4. 什么时候用 flow{}?

flow{} 适合的场景

// 1. 同步数据生成

fun generateNumbers(): Flow<Int> = flow {

    for (i in 1..10) {

        emit(i)

        delay(100)

    }

}

// 2. 简单的数据处理

fun processData(): Flow<String> = flow {

    val data = fetchData() // 同步操作

    emit(data)

}

// 3. 协程内的操作

fun apiCall(): Flow<String> = flow {

    val response = withContext(Dispatchers.IO) {

        api.getData() // 挂起函数

    }

    emit(response)

}

callbackFlow{} 适合的场景

// 1. 基于回调的 API

fun connectSSE(): Flow<SSEEvent> = callbackFlow {

    api.connect { event ->

        trySend(event) // 在回调中发送数据

    }

    awaitClose { api.disconnect() }

}

// 2. 事件监听

fun listenToEvents(): Flow<Event> = callbackFlow {

    eventListener.addListener { event ->

        trySend(event) // 在监听器中发送数据

    }

    awaitClose { eventListener.removeListener() }

}

// 3. 定时器

fun timer(): Flow<Unit> = callbackFlow {

    val timer = Timer()

    timer.scheduleAtFixedRate(object : TimerTask() {

        override fun run() {

            trySend(Unit) // 在定时器回调中发送数据

        }

    }, 0, 1000)

    awaitClose { timer.cancel() }

}

5. 实际例子对比

错误示例:用 flow{} 处理 SSE

// ❌ 这样不行

fun connectSSE(): Flow<SSEEvent> = flow {

    val request = Request.Builder()

        .url("https://api.example.com/events")

        .build()

    

    call = okHttpClient.newCall(request)

    call?.enqueue(object : Callback {

        override fun onResponse(call: Call, response: Response) {

            val body = response.body

            while (true) {

                val line = body?.source()?.readUtf8Line() ?: break

                if (line.startsWith("data:")) {

                    val data = line.substring(5).trim()

                    // ❌ 错误:不能在回调中使用 emit()

                    // emit(SSEEvent.Message("message", data))

                }

            }

        }

        

        override fun onFailure(call: Call, e: IOException) {

            // ❌ 错误:不能在回调中使用 emit()

            // emit(SSEEvent.Error(e))

        }

    })

    

    // ❌ 问题:flow 立即结束,不会等待 HTTP 响应

}

正确示例:用 callbackFlow{} 处理 SSE

// ✅ 这样正确

fun connectSSE(): Flow<SSEEvent> = callbackFlow {

    val request = Request.Builder()

        .url("https://api.example.com/events")

        .build()

    

    call = okHttpClient.newCall(request)

    call?.enqueue(object : Callback {

        override fun onResponse(call: Call, response: Response) {

            val body = response.body

            try {

                while (true) {

                    val line = body?.source()?.readUtf8Line() ?: break

                    if (line.startsWith("data:")) {

                        val data = line.substring(5).trim()

                        // ✅ 正确:在回调中使用 trySend()

                        trySend(SSEEvent.Message("message", data))

                    }

                }

                // ✅ 正确:发送关闭事件

                trySend(SSEEvent.Closed)

            } catch (e: Exception) {

                // ✅ 正确:发送错误事件

                trySend(SSEEvent.Error(e))

            } finally {

                response.close()

                close() // 关闭 Channel

            }

        }

        

        override fun onFailure(call: Call, e: IOException) {

            // ✅ 正确:发送错误事件

            trySend(SSEEvent.Error(e))

            close(e) // 关闭 Channel

        }

    })

    

    // ✅ 正确:等待 Channel 关闭

    awaitClose { 

        disconnect() 

    }

}

6. 总结

flow{} 的限制

  • ❌ 只能在协程中使用 emit()
  • ❌ 不能处理异步回调
  • ❌ 不能等待外部事件

callbackFlow{} 的优势

  • ✅ 可以在任何地方使用 trySend()
  • ✅ 专门处理基于回调的 API
  • ✅ 可以等待外部事件
  • ✅ 提供资源清理机制

选择原则

// 用 flow{} 当:

// - 数据生成是同步的

// - 在协程中操作

// - 简单的数据处理

// 用 callbackFlow{} 当:

// - 处理基于回调的 API

// - 需要等待外部事件

// - 需要资源清理

所以你的 SSE 代码必须使用 callbackFlow{},因为 HTTP 请求是基于回调的异步操作!

 

进一步分析

网络请求确实是异步的,而且确实可以用 flow{}。让我详细解释一下什么时候用 flow{},什么时候用 callbackFlow{}。

1. 网络请求的不同类型

类型一:挂起函数(可以用 flow{})


// 挂起函数 - 可以用 flow{}

suspend fun fetchData(): String {

    return withContext(Dispatchers.IO) {

        // 网络请求

        "数据"

    }

}

// 使用 flow{}

fun getDataFlow(): Flow<String> = flow {

    val data = fetchData() // 挂起函数

    emit(data)

}

类型二:回调函数(必须用 callbackFlow{})

 

// 回调函数 - 必须用 callbackFlow{}

fun fetchDataWithCallback(callback: (String) -> Unit) {

    // 异步网络请求

    api.getData { result ->

        callback(result)

    }

}

// 使用 callbackFlow{}

fun getDataFlow(): Flow<String> = callbackFlow {

    fetchDataWithCallback { result ->

        trySend(result) // 在回调中发送数据

    }

    awaitClose { /* 清理 */ }

}

2. 具体例子对比

例子一:Retrofit 挂起函数(用 flow{})

// Retrofit 接口

interface ApiService {

    @GET("data")

    suspend fun getData(): String // 挂起函数

}

// 使用 flow{}

fun getDataFlow(): Flow<String> = flow {

    val data = apiService.getData() // 挂起函数

    emit(data)

}

// 使用方式

getDataFlow().collect { data ->

    println("收到数据: $data")

}

例子二:OkHttp 回调(用 callbackFlow{})

// OkHttp 回调方式

fun getDataWithCallback(callback: (String) -> Unit) {

    val request = Request.Builder()

        .url("https://api.example.com/data")

        .build()

    

    client.newCall(request).enqueue(object : Callback {

        override fun onResponse(call: Call, response: Response) {

            val data = response.body?.string() ?: ""

            callback(data) // 回调

        }

        

        override fun onFailure(call: Call, e: IOException) {

            callback("错误: ${e.message}")

        }

    })

}

// 使用 callbackFlow{}

fun getDataFlow(): Flow<String> = callbackFlow {

    getDataWithCallback { result ->

        trySend(result) // 在回调中发送数据

    }

    awaitClose { /* 清理 */ }

}

3. 你的 SSE 代码分析

为什么 SSE 必须用 callbackFlow{}?

// SSE 使用的是 OkHttp 的回调方式

call?.enqueue(object : Callback {

    override fun onResponse(call: Call, response: Response) {

        // 在回调中处理数据

        while (true) {

            val line = body.source().readUtf8Line() ?: break

            if (line.startsWith("data:")) {

                val data = line.substring(5).trim()

                // 必须在回调中发送数据

                trySend(SSEEvent.Message("message", data))

            }

        }

    }

    

    override fun onFailure(call: Call, e: IOException) {

        // 在回调中处理错误

        trySend(SSEEvent.Error(e))

    }

})

 

如果 SSE 用挂起函数会怎样?

 

// 理论上可以这样,但 SSE 是长连接,不适合

interface SSEApi {

    @GET("events")

    suspend fun getEvents(): ResponseBody // 挂起函数

}

// 使用 flow{}(但这样不好)

fun getSSEFlow(): Flow<SSEEvent> = flow {

    val response = sseApi.getEvents() // 挂起函数

    val source = response.source()

    

    while (!source.exhausted()) {

        val line = source.readUtf8Line() ?: break

        if (line.startsWith("data:")) {

            val data = line.substring(5).trim()

            emit(SSEEvent.Message("message", data))

        }

    }

}

4. 实际项目中的例子

例子一:普通 API 请求(用 flow{})


// Retrofit 接口

interface UserApi {

    @GET("users")

    suspend fun getUsers(): List<User>

}

// 使用 flow{}

fun getUsersFlow(): Flow<List<User>> = flow {

    val users = userApi.getUsers() // 挂起函数

    emit(users)

}

// 使用方式

getUsersFlow().collect { users ->

    // 处理用户列表

}

例子二:WebSocket 连接(用 callbackFlow{})

 

// WebSocket 回调

fun connectWebSocket(): Flow<String> = callbackFlow {

    val webSocket = client.newWebSocket(request, object : WebSocketListener() {

        override fun onMessage(webSocket: WebSocket, text: String) {

            trySend(text) // 在回调中发送数据

        }

        

        override fun onFailure(webSocket: WebSocket, t: Throwable, response: Response?) {

            trySend("错误: ${t.message}")

        }

    })

    

    awaitClose { webSocket.close(1000, "正常关闭") }

}

例子三:定时器(用 callbackFlow{})

 

// 定时器回调

fun timerFlow(): Flow<Long> = callbackFlow {

    val timer = Timer()

    timer.scheduleAtFixedRate(object : TimerTask() {

        override fun run() {

            trySend(System.currentTimeMillis()) // 在回调中发送数据

        }

    }, 0, 1000)

    

    awaitClose { timer.cancel() }

}

5. 选择原则

用 flow{} 当:

 

// 1. 使用挂起函数

suspend fun fetchData(): String = api.getData()

fun getDataFlow(): Flow<String> = flow {

    val data = fetchData() // 挂起函数

    emit(data)

}

// 2. 在协程中操作

fun processDataFlow(): Flow<String> = flow {

    withContext(Dispatchers.IO) {

        val data = processData() // 挂起函数

        emit(data)

    }

}

// 3. 简单的数据转换

fun transformDataFlow(): Flow<String> = flow {

    val data = getData() // 挂起函数

    emit(data.uppercase())

}

用 callbackFlow{} 当:

 

// 1. 处理回调 API

fun connectSSE(): Flow<SSEEvent> = callbackFlow {

    api.connect { event ->

        trySend(event) // 在回调中发送数据

    }

    awaitClose { api.disconnect() }

}

// 2. 事件监听

fun listenToEvents(): Flow<Event> = callbackFlow {

    eventListener.addListener { event ->

        trySend(event) // 在监听器中发送数据

    }

    awaitClose { eventListener.removeListener() }

}

// 3. 长连接

fun connectWebSocket(): Flow<String> = callbackFlow {

    webSocket.connect { message ->

        trySend(message) // 在回调中发送数据

    }

    awaitClose { webSocket.close() }

}

6. 总结

关键区别

  • flow{}:用于挂起函数和协程操作
  • callbackFlow{}:用于回调函数和异步事件

你的 SSE 代码

 

// SSE 使用 OkHttp 的回调方式,所以必须用 callbackFlow{}

fun connect(url: String): Flow<SSEEvent> = callbackFlow {

    call?.enqueue(object : Callback {

        override fun onResponse(call: Call, response: Response) {

            // 在回调中处理数据

            trySend(SSEEvent.Message("message", data))

        }

    })

    awaitClose { disconnect() }

}

如果是挂起函数版本

 

// 如果 SSE 有挂起函数版本,可以用 flow{}

interface SSEApi {

    @GET("events")

    suspend fun getEvents(): ResponseBody

}

fun getSSEFlow(): Flow<SSEEvent> = flow {

    val response = sseApi.getEvents() // 挂起函数

    // 处理响应...

    emit(SSEEvent.Message("message", data))

}

所以你的理解是对的,网络请求确实可以用 flow{},但关键是要看网络请求的实现方式:

  • 挂起函数 → 用 flow{}
  • 回调函数 → 用 callbackFlow{}

你的 SSE 代码用的是 OkHttp 的回调方式,所以必须用 callbackFlow{}!

Logo

开源鸿蒙跨平台开发社区汇聚开发者与厂商,共建“一次开发,多端部署”的开源生态,致力于降低跨端开发门槛,推动万物智联创新。

更多推荐