Kotlin: Kotlin Flow详解
最后更新:2026-08-26
Flow 是协程的流式扩展——冷流按需发射、背压天然支持、操作符可组合。Charlie 用 构建实时订单处理管道。
1. 你将学到
- 冷流 vs 热流: 按需发射
- 中间操作: / / / /
- 终端操作: / / /
- 背压处理: / /
- Charlie 实战:订单事件流式处理管道
2. 一个架构师的真实故事
(1) 痛点:实时订单处理瓶颈
Charlie 的 OrderProcessor 需要实时处理每秒 5,000 笔订单事件。用 批量处理延迟高,用回调嵌套不可维护,用 RxJava 学习曲线陡峭。
(2) Flow 的解法
KOTLIN
// Flow: cold stream, on-demand emission, backpressure built-in
fun orderEventStream(): Flow`<OrderEvent>` = flow {
while (true) {
val event = pollOrderEvent() // Non-blocking poll
emit(event) // Emit to downstream
}
}
// Compose pipeline: filter -> transform -> collect
orderEventStream()
.filter { it.total > 100 }
.map { enrich(it) }
.collect { dispatch(it) }
Flow 是协程原生的流式 API——冷流按需发射、背压天然支持、操作符可组合、学习成本远低于 RxJava。
3. Flow 基础
(1) 创建 Flow
KOTLIN
// 1. flow builder
fun numberFlow(): Flow`<Int>` = flow {
for (i in 1..5) {
emit(i) // Emit value to collector
delay(100) // Simulate work
}
}
// 2. flowOf: from fixed values
val orderFlow = flowOf(
Order("ORD-001", 299.99),
Order("ORD-002", 1_500.00)
)
// 3. asFlow: from Iterable/Sequence
val orderList = listOf(Order("ORD-001", 299.99))
val orderFlow2 = orderList.asFlow()
// 4. channelFlow: for complex concurrent emission
fun concurrentFlow(): Flow`<Int>` = channelFlow {
launch { send(1) }
launch { send(2) }
}
(2) 冷流特性
KOTLIN
// Cold flow: each collector triggers new emission
val coldFlow = flow {
println("Emitting...") // This runs EACH time collect is called
emit(1)
emit(2)
}
coldFlow.collect { println("Collector 1: $it") } // Emits 1, 2
coldFlow.collect { println("Collector 2: $it") } // Emits 1, 2 AGAIN
(3) 冷流 vs 热流
| 维度 | 冷流(Flow) | 热流(SharedFlow/StateFlow) |
|---|---|---|
| 触发 | 收集时才发射 | 独立于收集者运行 |
| 数据 | 每个收集者收到完整数据 | 新收集者只收到后续数据 |
| 缓存 | 无 | 可配置 replay |
| 生命周期 | 随收集者 | 独立 |
| 类比 | 视频点播 | 直播流 |
4. 中间操作
中间操作是冷操作——不触发发射,只定义转换逻辑。
(1) map / filter / transform
KOTLIN
// map: transform each value
val totals = orderFlow.map { it.total }
// filter: keep matching values
val highValue = orderFlow.filter { it.total > 1_000 }
// transform: more flexible than map (can emit multiple values)
val expanded = orderFlow.transform { order ->
emit("Order: ${order.id}")
emit("Total: \$${order.total} USD")
}
(2) take / drop / debounce
KOTLIN
// take: take first N values
val first3 = orderFlow.take(3)
// drop: skip first N values
val after5 = orderFlow.drop(5)
// debounce: emit only after silence period
val debounced = searchFlow.debounce(300) // Wait 300ms of no new values
(3) 中间操作分类
| 类别 | 操作符 | 功能 |
|---|---|---|
| 转换 | / | 逐值转换 |
| 过滤 | / / | 条件过滤 |
| 展平 | / / | 流嵌套展平 |
| 截取 | / / | 取/跳前 N 个 |
| 时序 | / / | 时间控制 |
| 组合 | / | 多流组合 |
5. 终端操作
终端操作触发实际发射和收集。
(1) collect / toList / first
KOTLIN
// collect: consume each value
orderFlow.collect { order -> println(order.id) }
// toList: collect to list
val allOrders = orderFlow.toList()
// first: get first value
val firstOrder = orderFlow.first()
// single: expect exactly one value
val onlyOrder = orderFlow.single()
(2) 终端操作分类
| 操作 | 返回类型 | 说明 |
|---|---|---|
| Unit | 逐值消费 | |
| 收集为列表 | ||
| / | / | 取第一个值 |
| / | / | 取唯一值 |
| Int | 计数 | |
| 累积计算 |
6. 背压处理
(1) 背压问题
KOTLIN
// Producer is faster than consumer
val fastFlow = flow {
repeat(10) {
emit(it) // Fast: 1ms per emit
delay(1)
}
}
// Slow consumer causes backpressure
fastFlow.collect {
delay(100) // Slow: 100ms per process
println("Processed: $it")
}
// By default: producer waits for consumer (suspend)
(2) 背压策略
KOTLIN
// buffer: run producer and consumer in parallel
fastFlow.buffer(10) // Buffer up to 10 values
.collect { delay(100); println("Buffered: $it") }
// conflate: skip intermediate values, keep latest
fastFlow.conflate()
.collect { delay(100); println("Latest: $it") }
// collectLatest: cancel previous collection when new value arrives
fastFlow.collectLatest { value ->
delay(100)
println("Latest processed: $value")
}
(3) 背压策略对比
| 策略 | 行为 | 丢失数据 | 适用场景 |
|---|---|---|---|
| 默认 | 生产者等待消费者 | 不丢失 | 保证完整性 |
| 缓冲 N 个值 | 缓冲满时等待 | 短暂生产高峰 | |
| 跳过中间值,保留最新 | 丢失中间值 | 实时显示 | |
| 取消旧处理,处理最新 | 旧处理中断 | 搜索建议 |
7. Flow 冷流管道
flowchart LR
A[Source<br/>flow emit] --> B[filter<br/>total > 100]
B --> C[map<br/>enrich]
C --> D[debounce<br/>300ms]
D --> E[collect<br/>dispatch]
8. 完整示例:OrderProcessor 事件流
KOTLIN
// ============================================
// OrderProcessor - Event Stream with Flow
// Feature: Real-time order event processing
// ============================================
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
data class Order(val id: String, val total: Double, val status: String, val customer: String)
sealed class OrderEvent {
data class Created(val order: Order) : OrderEvent()
data class Updated(val orderId: String, val newStatus: String) : OrderEvent()
data class Cancelled(val orderId: String, val reason: String) : OrderEvent()
}
// Simulate event source
fun orderEventStream(): Flow`<OrderEvent>` = flow {
val events = listOf(
OrderEvent.Created(Order("ORD-001", 299.99, "PENDING", "Alice")),
OrderEvent.Created(Order("ORD-002", 15_000.00, "PENDING", "Bob")),
OrderEvent.Updated("ORD-001", "CONFIRMED"),
OrderEvent.Created(Order("ORD-003", 2_500.00, "PENDING", "Charlie")),
OrderEvent.Cancelled("ORD-003", "Out of stock"),
OrderEvent.Updated("ORD-002", "CONFIRMED"),
OrderEvent.Created(Order("ORD-004", 8_900.00, "PENDING", "Bob")),
OrderEvent.Updated("ORD-002", "SHIPPED")
)
events.forEach { event ->
emit(event)
delay(100)
}
}
// Enrich created orders
fun Order.enrich() = copy(status = "ENRICHED")
fun main() = runBlocking {
println("=== Order Event Stream ===")
// Pipeline: filter created -> enrich -> collect
orderEventStream()
.onStart { println("Stream started") }
.onEach { println(" Received: $it") }
.filter { it is OrderEvent.Created }
.map { (it as OrderEvent.Created).order }
.filter { it.total > 1_000 }
.map { it.enrich() }
.onEach { println(" >>> Enriched high-value: ${it.id} (\$${it.total} USD)") }
.onCompletion { println("Stream completed") }
.collect()
// Separate pipeline: status updates
println("\n=== Status Updates ===")
orderEventStream()
.filter { it is OrderEvent.Updated }
.map { it as OrderEvent.Updated }
.collect { event ->
println(" ${event.orderId} -> ${event.newStatus}")
}
// Aggregate: count events by type
println("\n=== Event Counts ===")
val counts = orderEventStream()
.groupBy { it::class.simpleName ?: "Unknown" }
.map { (type, flow) ->
type to flow.count()
}
.toList()
.associate { it }
counts.forEach { (type, count) ->
println(" $type: $count events")
}
}
输出:
TEXT
📖 仅展示
=== Order Event Stream ===
Stream started
Received: Created(order=Order(ORD-001, 299.99, PENDING, Alice))
Received: Created(order=Order(ORD-002, 15000.0, PENDING, Bob))
>>> Enriched high-value: ORD-002 ($15000.0 USD)
Received: Updated(orderId=ORD-001, newStatus=CONFIRMED)
Received: Created(order=Order(ORD-003, 2500.0, PENDING, Charlie))
>>> Enriched high-value: ORD-003 ($2500.0 USD)
Received: Cancelled(orderId=ORD-003, reason=Out of stock)
Received: Updated(orderId=ORD-002, newStatus=CONFIRMED)
Received: Created(order=Order(ORD-004, 8900.0, PENDING, Bob))
>>> Enriched high-value: ORD-004 ($8900.0 USD)
Received: Updated(orderId=ORD-002, newStatus=SHIPPED)
Stream completed
=== Status Updates ===
ORD-001 -> CONFIRMED
ORD-002 -> CONFIRMED
ORD-002 -> SHIPPED
=== Event Counts ===
Created: 4 events
Updated: 3 events
Cancelled: 1 events
❓ 常见问题
Q Flow 和 Sequence 有什么区别?
A Flow 是异步的(支持 suspend),Sequence 是同步的。Flow 适合异步数据源(网络/数据库),Sequence 适合同步大数据处理。
Q Flow 和 RxJava 的 Observable 有什么区别?
A Flow 是协程原生的、冷流优先的、学习曲线平缓的。RxJava 更成熟但更复杂。新项目推荐 Flow,已有 RxJava 的项目可共存。
Q 冷流什么时候变热流?
A 用 将冷流转为 SharedFlow,用 转为 StateFlow。这让多个收集者共享同一个数据源。
Q flatMapConcat、flatMapMerge、flatMapLatest 怎么选?
A 顺序处理(默认), 并发处理(提高吞吐), 取消旧流(实时场景如搜索)。
Q Flow 的背压默认怎么处理?
A 默认是"生产者等待消费者"——emit 是 suspend 函数,消费者慢时生产者自动挂起。这是最安全的默认行为。
Q collect 是挂起函数吗?
A 是的。 是 suspend 函数,只能在协程或 suspend 函数中调用。它会挂起直到 Flow 完成。
📖 小节
- Flow 是冷流:每个收集者触发独立发射,不共享数据
- 中间操作(map/filter/transform)是冷的,不触发发射
- 终端操作(collect/toList/first)触发实际数据流
- 背压默认策略:生产者等待消费者,可通过 buffer/conflate/collectLatest 调整
- Flow 操作符与集合操作符语法一致,学习成本低
- // 用于流的生命周期钩子
📝 作业
- 基础题(难度⭐):用 创建一个发射 1-10 的 Flow,用 打印偶数。提示:
- 进阶题(难度⭐⭐):用 Flow 模拟订单事件流,用 + + 处理金额 > 1000 的订单。提示:参考第 8 节
- 挑战题(难度⭐⭐⭐):实现带 的搜索流:模拟用户输入,300ms 内无新输入才触发搜索。提示: