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

6.0 KiB
Raw Blame History

tags, create time
tags create time
MQ
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 配合是最常见的组合:

  • 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 适合重量级场景(日均十亿级、复杂窗口逻辑、需要强一致性保证)。

关联笔记