6.0 KiB
tags, create time
| tags | create time | |
|---|---|---|
|
2026-05-24 19:52 |
MQ 与流处理
概述
消息队列处理的是离散的消息,消费即删;流处理平台处理的是连续的数据流,持久保留、可回溯。本文对比两者的差异,深入流处理的核心概念(事件时间、窗口、Watermark),并介绍 Kafka Streams、Flink、Spark Streaming 三大流处理方案的特点与适用场景。
正文
消息队列 vs 流处理平台
很多人把 Kafka 叫"消息队列",但它其实更接近一个流处理平台。两者的本质区别在于:
| 维度 | 消息队列 | 流处理平台 |
|---|---|---|
| 数据模型 | 离散消息,消费即删 | 连续流,持久保留 |
| 消费语义 | 一条消息只被消费一次 | 同一数据可被多次回溯 |
| 处理方式 | 单条处理 | 窗口聚合、流式计算 |
| 典型代表 | RabbitMQ, RocketMQ | Kafka, Pulsar, Flink |
[!question] Kafka 到底是消息队列还是流处理平台?这个争论有意义吗?
其实没有意义。Kafka 是一个分布式日志系统,你可以把它当消息队列用,也可以当流处理平台用。关键不在于它"是什么",而在于你怎么用它。
流处理的核心概念
事件时间 vs 处理时间
这是流处理中最容易混淆的概念:
- 事件时间(Event Time):事件实际发生的时间,由消息自身携带
- 处理时间(Processing Time):事件被流处理引擎处理的时间
两者可能相差毫秒,也可能相差小时(网络延迟、积压、重放)。正确使用事件时间是保证结果准确性的前提。
窗口(Window)
流数据是无限的,但业务计算需要有限的数据集。窗口把无限流切分成有限的"片段":
- 滚动窗口(Tumbling):固定大小,不重叠。每 5 分钟一个窗口
- 滑动窗口(Sliding):固定大小,可重叠。窗口大小 10 分钟,每 1 分钟滑动一次
- 会话窗口(Session):按活跃度切分,超时即关闭窗口
Watermark
Watermark 解决的是"迟到数据"问题。它是一个时间戳,表示"在这个时间之前的数据,我认为已经到齐了"。
当 Watermark 推过窗口的结束时间,窗口就会触发计算并关闭。但如果数据迟到超过 Watermark,就只能靠 允许迟到(Allowed Lateness) 或 侧输出(Side Output) 来补救。
Kafka Streams
Kafka Streams 是一个轻量级流处理库,不需要独立集群,直接嵌入应用中运行。
核心抽象:
- KStream:无界的、追加式的记录流(类似日志)
- KTable:可更新的 changelog 流(类似数据库表)
// Kafka Streams 的 DSL 思路(伪代码)
// 从 topic 读取流 -> 按 key 分组 -> 聚合 -> 写回 topic
stream := builder.Stream("input-topic")
stream.GroupByKey().
WindowedBy(TimeWindows.OfSize(5 * time.Minute)).
Count().
ToStream().
To("output-topic")
Kafka Streams 的优势在于运维简单——不需要 Flink 那样的独立集群,应用启动就自动处理,应用停止就自动释放资源。
Flink + Kafka
Flink 是真正的分布式流处理引擎,与 Kafka 配合是最常见的组合:
- Source/Sink Connector:Flink 原生支持 Kafka 作为数据源和输出目标
- Exactly-Once 保障:Flink 的 checkpoint 机制 + Kafka 的事务消费,端到端 exactly-once
- 状态后端:RocksDB 状态后端支持超大状态(TB 级别),适合复杂聚合
Flink 的核心优势是 低延迟 + 高吞吐 + 强一致性。适合对延迟和正确性要求极高的场景。
Spark Streaming + Kafka
Spark Streaming 使用微批处理模型:把流数据切成小批次(通常秒级),每个批次用 Spark 引擎处理。
严格来说这不是"真正的流处理",但在很多场景下足够用。优势是复用了 Spark 生态(ML、SQL、GraphX),适合批流一体的需求。Structured Streaming 在此基础上提供了更接近真正流处理的 API。
流处理适用场景
- 实时聚合:实时 PV/UV 统计、实时 GMV 计算
- 实时风控:检测异常登录、欺诈交易,毫秒级响应
- 实时推荐:根据用户实时行为更新推荐模型
- IoT 数据处理:传感器数据实时清洗、聚合、告警
graph LR
MQ["消息队列"] -->|"数据流"| SP["流处理器"]
SP -->|"聚合结果"| DB["数据库"]
SP -->|"实时告警"| A["告警系统"]
SP -->|"写回"| MQ
style MQ fill:#4CAF50,color:#fff
style SP fill:#2196F3,color:#fff
style DB fill:#FF9800,color:#fff
style A fill:#F44336,color:#fff
Go 代码:简化版窗口聚合
// 简化版窗口聚合逻辑
type Event struct {
Timestamp time.Time
Value float64
}
func windowAggregate(events <-chan Event, windowSize time.Duration) {
window := make([]Event, 0)
ticker := time.NewTicker(windowSize)
defer ticker.Stop()
for {
select {
case e := <-events:
window = append(window, e)
case <-ticker.C:
// 窗口关闭,触发聚合计算
sum := 0.0
for _, e := range window {
sum += e.Value
}
avg := sum / float64(len(window))
fmt.Printf("Window closed: count=%d, avg=%.2f\n", len(window), avg)
window = window[:0] // 清空窗口
}
}
}
这段代码展示了窗口聚合的核心思想:在一个时间窗口内收集事件,窗口关闭时触发计算。真实场景中还需要处理迟到数据、状态持久化、checkpoint 等问题,但基本思路是一样的。
[!question] Kafka Streams 和 Flink 都能做流处理,什么时候选哪个?
简单说:Kafka Streams 适合轻量级场景(日均百万级、逻辑简单、不想运维额外集群),Flink 适合重量级场景(日均十亿级、复杂窗口逻辑、需要强一致性保证)。