Kotlin: Kotlin Flow详解

最后更新:2026-08-26

Flow 是协程的流式扩展——冷流按需发射、背压天然支持、操作符可组合。Charlie 用 构建实时订单处理管道。

1. 你将学到


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 冷流管道

100%
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 完成。

📖 小节


📝 作业

  1. 基础题(难度⭐):用 创建一个发射 1-10 的 Flow,用 打印偶数。提示:
  2. 进阶题(难度⭐⭐):用 Flow 模拟订单事件流,用 + + 处理金额 > 1000 的订单。提示:参考第 8 节
  3. 挑战题(难度⭐⭐⭐):实现带 的搜索流:模拟用户输入,300ms 内无新输入才触发搜索。提示:

← 上一课 | 下一课 →

Web-Tutorial.com

Web-Tutorial 技术团队

由多位开发者共同维护的编程教程平台。每篇教程由对应领域的开发者编写和审核,确保内容准确可靠。如发现任何问题,欢迎向我们反馈。

100%

🙏 帮我们做得更好

我们是刚上线的编程教程站,几个人的小团队,精力有限。页面虽经检查,难免还有疏漏——链接失效、排版错乱、内容有误、语言生硬……

如果您发现了,麻烦告诉我们,我们会在收到反馈后第一时间进行修复,再次感谢您的光临 🙏