flow{} 和 callbackFlow{} 的区别
·
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{}!
更多推荐


所有评论(0)