Files
cs-note/hhs/MQ/09-流处理与事件驱动/33-MQ-与流处理.md
T
2026-05-24 20:51:06 +08:00

157 lines
6.0 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
---
tags:
- MQ
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 流(类似数据库表)
```go
// 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 数据处理**:传感器数据实时清洗、聚合、告警
```mermaid
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 代码:简化版窗口聚合
```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 适合重量级场景**(日均十亿级、复杂窗口逻辑、需要强一致性保证)。
## 关联笔记
- [[31-MQ-背压与流控]]
- [[32-MQ-请求-回复模式]]
- [[34-事件驱动架构-EDA]]