返回文章列表
技术2026年9月18日17 分钟阅读

流处理系统:Flink与Spark Streaming的实时计算原理

流处理系统:Flink与Spark Streaming的实时计算原理

一、从批处理到流处理:计算范式的演进

1.1 数据处理的两种模式

批处理(Batch Processing):

对有限、静态的数据集进行计算,强调高吞吐。

批处理特征:
├── 输入:有限数据集(如一天日志)
├── 处理:全量读取 → 计算 → 输出
├── 延迟:分钟级到小时级
├── 代表:Hadoop MapReduce、Spark SQL
└── 适用:离线分析、报表生成

流处理(Stream Processing):

对无限、动态的数据流进行实时计算,强调低延迟。

流处理特征:
├── 输入:无限数据流(如实时点击流)
├── 处理:来一条处理一条
├── 延迟:毫秒级到秒级
├── 代表:Flink、Spark Streaming、Kafka Streams
└── 适用:实时监控、风控、推荐

1.2 Lambda架构与Kappa架构

Lambda架构(Nathan Marz, 2011):

Lambda 架构:

         ┌─────────────┐
         │  数据源      │
         └──────┬──────┘
                │
       ┌────────┴────────┐
       ▼                 ▼
┌─────────────┐   ┌─────────────┐
│  批处理层    │   │  速度层      │
│  (Batch)    │   │  (Speed)    │
│  全量计算    │   │  增量计算    │
│  高准确      │   │  低延迟      │
└──────┬──────┘   └──────┬──────┘
       │                 │
       ▼                 ▼
┌─────────────┐   ┌─────────────┐
│  批处理视图  │   │  实时视图    │
│  (准确)      │   │  (近似)      │
└──────┬──────┘   └──────┬──────┘
       │                 │
       └────────┬────────┘
                ▼
         ┌─────────────┐
         │  服务层      │
         │  (合并视图)  │
         └─────────────┘

问题:
- 维护两套代码(批处理 + 流处理)
- 逻辑不一致风险
- 系统复杂度高

Kappa架构(Jay Kreps, 2014):

Kappa 架构:

         ┌─────────────┐
         │  数据源      │
         └──────┬──────┘
                │
                ▼
         ┌─────────────┐
         │  消息队列    │
         │  (Kafka)    │
         └──────┬──────┘
                │
                ▼
         ┌─────────────┐
         │  流处理层    │
         │  (唯一处理)  │
         │  统一代码    │
         └──────┬──────┘
                │
                ▼
         ┌─────────────┐
         │  服务层      │
         └─────────────┘

改进:
- 只有一套代码(流处理)
- 重放机制实现批处理
- 系统简化

1.3 流处理的核心挑战

1. 无限数据:

  • 内存无法容纳全部数据
  • 需要窗口(Window)机制

2. 时间语义:

  • 事件时间(Event Time)vs 处理时间(Processing Time)
  • 乱序数据与水位线(Watermark)

3. 状态管理:

  • 算子状态的持久化
  • 故障恢复与一致性

4. exactly-once 语义:

  • 不丢失数据
  • 不重复处理

二、时间语义与窗口机制

2.1 三种时间概念

时间类型:

事件时间(Event Time):
├── 数据产生的时间
├── 嵌入在数据记录中
├── 最准确,但可能有乱序
└── 需要处理延迟数据

摄取时间(Ingestion Time):
├── 数据进入系统的时间
├── 由消息队列设置
├── 介于事件时间和处理时间之间

处理时间(Processing Time):
├── 算子处理数据的时间
├── 最简单,延迟最低
└── 不确定性最高

乱序数据示例:

实际发生顺序:
T=1s  事件 A (用户点击)
T=2s  事件 B (用户点击)
T=3s  事件 C (用户点击)

到达系统顺序(网络延迟):
T=1s  事件 A
T=5s  事件 C  (延迟 2s)
T=6s  事件 B  (延迟 4s)

问题:
- 如果按处理时间统计,T=1-3s 窗口只包含 A
- 但 B 和 C 实际也属于这个时间段

2.2 水位线(Watermark)

水位线机制:

一种表示"事件时间进展"的特殊标记,用于处理乱序数据。

水位线定义:
Watermark(t) = 当前最大事件时间 - 允许延迟

含义:
- 所有事件时间 < t 的数据都已到达
- 可以触发窗口计算
- 延迟数据(> t)可以丢弃或特殊处理

水位线传播:

python
class WatermarkStrategy:
    """
    水位线生成策略
    """
    
    @staticmethod
    def for_bounded_out_of_orderness(delay: timedelta):
        """
        有界乱序水位线
        适用于最大延迟已知的场景
        """
        return BoundedOutOfOrdernessGenerator(delay)
    
    @staticmethod
    def for_monotonous_timestamps():
        """
        单调递增水位线
        适用于无乱序场景
        """
        return AscendingTimestampGenerator()


class BoundedOutOfOrdernessGenerator:
    """
    有界乱序水位线生成器
    """
    
    def __init__(self, max_out_of_orderness: timedelta):
        self.max_delay = max_out_of_orderness
        self.max_timestamp = datetime.min
    
    def extract_timestamp(self, event) -> datetime:
        """提取事件时间戳"""
        timestamp = event['timestamp']
        self.max_timestamp = max(self.max_timestamp, timestamp)
        return timestamp
    
    def get_current_watermark(self) -> datetime:
        """生成当前水位线"""
        return self.max_timestamp - self.max_delay

水位线可视化:

时间轴:
0s    1s    2s    3s    4s    5s    6s
│     │     │     │     │     │     │
├─────┼─────┼─────┼─────┼─────┼─────┤
│  A  │     │     │  C  │  B  │     │  ← 事件到达
│     │     │     │     │     │     │
└─────┴─────┴─────┴─────┴─────┴─────┘

水位线进展(max_delay=2s):
T=1s: Watermark = 1s - 2s = -1s (不触发)
T=5s: Watermark = 3s - 2s = 1s  (触发 [0-1s] 窗口)
T=6s: Watermark = 3s - 2s = 1s  (等待更多数据)

2.3 窗口类型

滚动窗口(Tumbling Window):

固定大小,不重叠

时间:  0-5s  5-10s  10-15s  15-20s
       ┌─────┐┌─────┐┌─────┐┌─────┐
事件:  │A B C││D E  ││F G H││I J  │
       └─────┘└─────┘└─────┘└─────┘

滑动窗口(Sliding Window):

固定大小,可重叠

窗口大小=10s,滑动步长=5s

时间:  0-10s   5-15s   10-20s
       ┌───────┐
       │A B C D│
       └───┬───┘
           │┌───────┐
           ││C D E F│
           │└───┬───┘
           │    │┌───────┐
           │    ││E F G H│
           │    │└───────┘

会话窗口(Session Window):

动态大小,由活动间隙决定

间隙(Gap)= 5s

事件:  A    B         C    D    E              F
时间:  1s   2s        10s  11s  12s            25s
       ├────┤         ├────┴────┤              │
       │窗口1│         │   窗口2  │              │窗口3
       └────┘         └─────────┘              │
                                                │
       A,B 间隔 < 5s → 同一窗口                │
       B,C 间隔 > 5s → 新窗口                  │
       C,D,E 连续 → 同一窗口                   │
       E,F 间隔 > 5s → 新窗口                  │

全局窗口(Global Window):

只有一个窗口,所有数据都在其中
需要自定义触发器(Trigger)决定何时计算

适用场景:
- 自定义聚合逻辑
- 复杂事件处理(CEP)

三、Flink:真正的流处理引擎

3.1 Flink架构概览

Flink 架构:

┌─────────────────────────────────────────┐
│           Flink Client                  │
│  (提交作业,生成 JobGraph)               │
└──────────────────┬──────────────────────┘
                   │
                   ▼
┌─────────────────────────────────────────┐
│      JobManager (主节点)                 │
│  ┌─────────────┐  ┌─────────────────┐   │
│  │ Dispatcher  │  │ ResourceManager │   │
│  │ (接收作业)   │  │ (资源管理)       │   │
│  └─────────────┘  └─────────────────┘   │
│  ┌─────────────────────────────────────┐│
│  │      JobMaster (作业管理)            ││
│  │  - 调度执行图 (ExecutionGraph)       ││
│  │  - 协调 Checkpoint                  ││
│  │  - 故障恢复                         ││
│  └─────────────────────────────────────┘│
└──────────────────┬──────────────────────┘
                   │
                   ▼
┌─────────────────────────────────────────┐
│      TaskManager (工作节点) × N          │
│  ┌─────────────────────────────────────┐│
│  │  Slot 1    │ Slot 2    │ Slot 3    ││
│  │  (任务槽)   │ (任务槽)   │ (任务槽)   ││
│  │  ┌──────┐  │  ┌──────┐  │  ┌──────┐ ││
│  │  │ Task │  │  │ Task │  │  │ Task │ ││
│  │  └──────┘  │  └──────┘  │  └──────┘ ││
│  └────────────┴────────────┴───────────┘│
└─────────────────────────────────────────┘

3.2 DataStream API

java
// Flink DataStream 示例
StreamExecutionEnvironment env = 
    StreamExecutionEnvironment.getExecutionEnvironment();

// 设置事件时间和水位线
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);

DataStream<ClickEvent> clicks = env
    .addSource(new KafkaSource<>("click-stream"))
    .assignTimestampsAndWatermarks(
        WatermarkStrategy
            .<ClickEvent>forBoundedOutOfOrderness(
                Duration.ofSeconds(5)
            )
            .withTimestampAssigner(
                (event, timestamp) -> event.getEventTime()
            )
    );

// 窗口聚合
DataStream<PageViewCount> pageViews = clicks
    .keyBy(ClickEvent::getPageId)
    .window(TumblingEventTimeWindows.of(Time.minutes(1)))
    .aggregate(new CountAggregate());

// 执行
env.execute("Page View Analytics");

3.3 Flink的状态管理

状态类型:

java
// Keyed State(键控状态)
// 每个 key 有独立的状态
ValueState<Integer> countState;
ListState<Event> eventListState;
MapState<String, Integer> mapState;
ReducingState<Event> reducingState;
AggregatingState<Event, Integer> aggState;

// Operator State(算子状态)
// 每个算子实例共享
ListState<Integer> operatorState;
BroadcastState<String, Rule> broadcastState;

状态后端:

MemoryStateBackend:
├── 状态存储在 JVM 堆内存
├── Checkpoint 存储在 JobManager 内存
├── 适用:本地测试、小状态
└── 限制:状态大小受内存限制

FsStateBackend:
├── 状态存储在 JVM 堆内存
├── Checkpoint 存储在文件系统
├── 适用:大状态、生产环境
└── 限制:状态大小仍受内存限制

RocksDBStateBackend:
├── 状态存储在 RocksDB(本地磁盘)
├── Checkpoint 存储在文件系统
├── 适用:超大状态、增量 Checkpoint
└── 优势:状态大小不受内存限制

3.4 Checkpoint机制

分布式快照(Chandy-Lamport算法):

Checkpoint 流程:

1. JobManager 向所有 Source 发送 Checkpoint Barrier
   
   Source ──> Barrier ──> 下游算子

2. Source 快照自己的状态,向下游传播 Barrier

3. 算子收到 Barrier 后:
   a. 将当前处理的数据缓冲
   b. 快照自己的状态
   c. 向下游传播 Barrier

4. 所有算子完成快照后,向 JobManager 确认

5. JobManager 收到所有确认后,Checkpoint 完成

Barrier对齐(Alignment):

多输入流的 Barrier 对齐:

算子有两个输入流:
Stream A:  A1, A2, [Barrier], A3, A4
Stream B:  B1, B2, B3, [Barrier], B4

处理过程:
1. 正常处理 A1, A2, B1, B2, B3
2. 收到 Stream A 的 Barrier
3. 暂停处理 Stream A,缓存 A3, A4
4. 继续处理 Stream B
5. 收到 Stream B 的 Barrier
6. 快照状态
7. 继续处理缓存的 A3, A4 和 B4

非对齐 Checkpoint(Flink 1.11+):
- 不暂停处理
- 将"在途数据"也作为状态的一部分
- 降低延迟,但状态更大

3.5 Exactly-Once语义

端到端 Exactly-Once:

三个层面的保证:

1. Flink 内部(Checkpoint)
   └── 分布式快照保证状态一致性

2. Source(可重放)
   └── Kafka:offset 作为状态
   └── 文件:读取位置作为状态

3. Sink(幂等或事务)
   ├── 幂等写入:Elasticsearch、Redis
   └── 事务写入:Kafka、JDBC

Kafka 事务写入示例:
- 开启事务
- 消费数据 → 处理 → 写入 Kafka
- Checkpoint 时提交事务
- 失败回滚到上一次 Checkpoint

四、Spark Streaming:微批处理架构

4.1 DStream与微批

Spark Streaming 核心思想:

将流处理转化为一系列小批处理(Micro-batch)。

微批处理:

实时流: ──────────────────────────────────────►
         │    │    │    │    │    │    │
         ▼    ▼    ▼    ▼    ▼    ▼    ▼
Batch:  [B1] [B2] [B3] [B4] [B5] [B6] [B7]
         │    │    │    │    │    │    │
         ▼    ▼    ▼    ▼    ▼    ▼    ▼
处理:   P1   P2   P3   P4   P5   P6   P7
         │    │    │    │    │    │    │
         └────┴────┴────┴────┴────┴────┘
                    输出

特点:
- 批大小:通常 100ms - 几秒
- 延迟:秒级
- 吞吐:高(利用 Spark SQL 优化)

4.2 Structured Streaming

Spark 2.0 引入 Structured Streaming,统一批处理和流处理 API:

scala
// Structured Streaming 示例
val spark = SparkSession.builder()
  .appName("StructuredStreaming")
  .getOrCreate()

import spark.implicits._

// 读取流数据
val lines = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "localhost:9092")
  .option("subscribe", "events")
  .load()
  .selectExpr("CAST(value AS STRING)")

// 处理(与批处理相同 API)
val wordCounts = lines
  .flatMap(_.split(" "))
  .groupBy("value")
  .count()

// 输出
val query = wordCounts.writeStream
  .outputMode("complete")  // complete/update/append
  .format("console")
  .start()

query.awaitTermination()

输出模式:

模式 描述 适用场景
Complete 输出完整结果表 聚合操作
Update 只输出更新的行 有键的聚合
Append 只输出新增的行 无聚合、窗口

4.3 Spark Streaming的状态管理

UpdateStateByKey:

scala
// 维护全局状态
val stateDStream = events
  .map(event => (event.userId, event.value))
  .updateStateByKey[Int](
    // 更新函数
    (newValues: Seq[Int], state: Option[Int]) => {
      val currentSum = state.getOrElse(0)
      val newSum = newValues.sum + currentSum
      Some(newSum)
    },
    // 分区器
    new HashPartitioner(sc.defaultParallelism)
  )

MapGroupsWithState:

scala
// 更灵活的状态管理
val sessionized = events
  .groupByKey(event => event.sessionId)
  .mapGroupsWithState[SessionState, SessionResult](
    GroupStateTimeout.NoTimeout()
  ) { (sessionId, events, state) =>
    
    val currentState = state.getOption.getOrElse(SessionState.empty)
    val newState = events.foldLeft(currentState)(_ + _)
    
    if (newState.isExpired) {
      state.remove()
      SessionResult(sessionId, newState, expired = true)
    } else {
      state.update(newState)
      SessionResult(sessionId, newState, expired = false)
    }
  }

4.4 Checkpoint与WAL

Checkpoint机制:

Spark Streaming Checkpoint:

1. 元数据 Checkpoint:
   ├── 配置信息
   ├── DStream 操作链
   └── 未完成的批次

2. 数据 Checkpoint:
   └── 有状态操作的状态数据
   └── 存储在 HDFS/S3

恢复流程:
1. 从 Checkpoint 读取元数据
2. 重建 DStream 图
3. 恢复状态数据
4. 继续处理

WAL(Write Ahead Log):

接收数据时的 WAL:

1. 从 Kafka 接收数据
2. 先写入 WAL(HDFS)
3. 再处理数据
4. 处理完成后删除 WAL

作用:
- Driver 故障时,从 WAL 恢复数据
- 保证数据不丢失
- 代价:写入延迟

五、Flink vs Spark Streaming 深度对比

5.1 架构对比

特性 Flink Spark Streaming
处理模型 原生流处理 微批处理
延迟 毫秒级 秒级
吞吐 高 很高(批处理优化)
状态管理 内置,多种后端 需借助外部系统
Exactly-Once 原生支持 需额外配置
事件时间支持 完善 较完善
SQL支持 Flink SQL Spark SQL(更成熟)
生态集成 Kafka、ES、HBase 更丰富的连接器

5.2 适用场景

选择 Flink:

  • 延迟要求 < 1秒
  • 复杂事件处理(CEP)
  • 需要精确的事件时间处理
  • 大状态场景(> 100GB)

选择 Spark Streaming:

  • 延迟要求 > 1秒可接受
  • 需要与 Spark 生态深度集成
  • 大量 SQL/ML 处理
  • 已有 Spark 团队

5.3 性能对比

基准测试(Yahoo Streaming Benchmark):

指标 Flink Spark Streaming Storm
吞吐 (events/s) 10M+ 20M+ 5M
延迟 (p99) 50ms 500ms 100ms
CPU 使用率 中等 较高 高
内存使用 中等 较高 低

六、流处理系统实践

6.1 背压处理

背压(Backpressure):

当消费速度 < 生产速度时,系统需要减速上游生产,避免 OOM。

Flink 背压机制:

下游 Task 繁忙 → 缓冲区满 → 停止读取上游
    ↑                              │
    └──────────────────────────────┘

自动调节:
- 降低 Source 消费速度
- 动态调整 Checkpoint 间隔
- 自动扩缩容(K8s 集成)

6.2 监控与调优

关键指标:

python
STREAMING_METRICS = {
    # 延迟指标
    'records_lag': '消费延迟(Kafka)',
    'watermark_delay': '水位线延迟',
    'processing_delay': '处理延迟',
    
    # 吞吐指标
    'records_in_per_second': '输入吞吐',
    'records_out_per_second': '输出吞吐',
    
    # 状态指标
    'state_size': '状态大小',
    'checkpoint_duration': 'Checkpoint 耗时',
    'checkpoint_size': 'Checkpoint 大小',
    
    # 资源指标
    'cpu_usage': 'CPU 使用率',
    'memory_usage': '内存使用率',
    'gc_time': 'GC 时间',
}

调优建议:

1. 并行度设置
   └── 通常与 Kafka 分区数相同
   
2. 缓冲区大小
   └── 增加缓冲区可减少网络往返
   
3. Checkpoint 间隔
   └── 权衡:恢复速度 vs 性能影响
   
4. 状态后端选择
   └── 大状态 → RocksDB
   └── 小状态 → Heap

6.3 实际案例:实时风控系统

架构:

用户行为 ──> Kafka ──> Flink CEP ──> 规则引擎 ──> 决策
                           │
                           ▼
                    历史特征(Redis)
                           ▲
                           │
                    特征计算(Flink)

Flink 作业:
1. 实时特征计算
   - 滑动窗口统计
   - 会话特征提取
   
2. CEP 复杂事件检测
   - 异常登录模式
   - 高频交易检测
   
3. 规则评分
   - 多维度风险评分
   - 实时决策输出

性能:
- 延迟:P99 < 200ms
- 吞吐:100K+ events/s
- 可用性:99.99%

结语

流处理系统代表了大数据计算从"离线分析"到"实时决策"的范式转变。Flink 和 Spark Streaming 代表了两种架构哲学:

Flink:

"流处理优先,批处理是流的特例"

  • 原生流处理架构
  • 低延迟、高准确
  • 适合实时性要求高的场景

Spark Streaming:

"批处理统一流处理"

  • 微批处理架构
  • 高吞吐、生态丰富
  • 适合已有 Spark 生态的场景

理解流处理系统,不仅是掌握技术工具,更是理解实时计算的本质:在无限的数据流中,如何及时、准确地提取价值。


参考资源

经典论文:

  1. Akidau, T., et al. (2013). "MillWheel: Fault-Tolerant Stream Processing at Internet Scale". VLDB.
  2. Carbone, P., et al. (2015). "Apache Flink: Stream and Batch Processing in a Single Engine". IEEE Data Engineering Bulletin.
  3. Zaharia, M., et al. (2013). "Discretized Streams: Fault-Tolerant Streaming Computation at Scale". SOSP.

官方文档: 4. Flink 官方文档:https://nightlies.apache.org/flink/flink-docs-stable/ 5. Spark Streaming 指南:https://spark.apache.org/docs/latest/streaming-programming-guide.html 6. Kafka Streams 文档:https://kafka.apache.org/documentation/streams/

实践资源: 7. 《Streaming Systems》(Tyler Akidau) 8. Flink 中文社区:https://flink-learning.org.cn/


创建时间:2026年04月11日
更新时间:2026年04月11日