Structured Streaming 核心
概述
Structured Streaming 把实时流当作无限表(Unbounded Table),用 DataFrame/Dataset API 统一流批处理。核心能力:微批/连续两种执行模式、Event-time 处理、Watermark 水位线与状态管理。本文逐一讲透。
一、无限表模型
1.1 核心思想
输入流 = 无限追加的表
查询 = 对表的持续查询(每批数据触发一次)
输出 = 增量更新到结果表| 概念 | 说明 |
|---|---|
| 输入表 | 数据到达即追加一行 |
| 查询 | 静态 SQL/DataFrame 定义 |
| 结果表 | 每批更新(追加/更新/完整) |
| 触发器 | 控制批次执行节奏 |
1.2 代码形态
scala
val df = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "kafka:9092")
.option("subscribe", "orders")
.load()
.selectExpr("CAST(value AS STRING) AS json")
val result = df.writeStream
.outputMode("append")
.format("parquet")
.option("path", "hdfs:///out/orders")
.option("checkpointLocation", "hdfs:///ckpt/orders")
.start()二、两种执行模式
2.1 微批处理(默认)
| 特点 | 说明 |
|---|---|
| 触发执行 | 按触发器间隔执行批次 |
| 容错 | 复用 Spark 批处理容错 |
| 延迟 | 秒级(默认 100ms 探测) |
2.2 连续处理(实验)
| 特点 | 说明 |
|---|---|
| 执行方式 | 常驻算子持续处理 |
| 延迟 | 毫秒级 |
| 限制 | 仅部分算子支持(map/filter 等) |
| 对比 | 微批 | 连续 |
|---|---|---|
| 延迟 | 秒级 | 毫秒级 |
| 算子支持 | 全量 | 有限 |
| 稳定性 | 生产成熟 | 实验性质 |
三、Event-time 处理
3.1 为什么需要 Event-time
| 时间类型 | 说明 |
|---|---|
| Event-time | 数据产生时间(业务字段) |
| Processing-time | 到达处理时间 |
流数据乱序/延迟到达,按 Processing-time 会算错,必须按 Event-time 聚合。
3.2 指定 Event-time
scala
val withEventTime = df
.withColumn("event_time", to_timestamp(col("ts")))
val windowed = withEventTime
.withWatermark("event_time", "10 minutes")
.groupBy(
window(col("event_time"), "10 minutes"),
col("user_id")
)
.count()| 关键点 | 说明 |
|---|---|
| 时间列 | 数据中的业务时间戳 |
| 窗口函数 | window 按事件时间切窗 |
| Watermark | 容忍延迟的上限 |
四、Watermark 水位线
4.1 原理
Watermark = 已处理数据的最大 Event-time - 容忍延迟
例:watermark 10 分钟
数据最大 event-time = 12:00
水位线 = 11:50
→ event-time > 11:50 的数据参与聚合
→ event-time < 11:50 的迟到数据被丢弃| 作用 | 说明 |
|---|---|
| 定义"迟到" | 超过水位线的数据视为迟到 |
| 触发窗口关闭 | 窗口过期即输出最终结果 |
| 清理状态 | 删除过期窗口状态 |
4.2 迟到的处理
| 情况 | 行为 |
|---|---|
| event-time > 水位线 | 正常聚合(补入对应窗口) |
| 已关闭窗口的迟到数据 | 被丢弃(默认) |
| 需要全部保留 | 用追加模式的过滤器或侧输出 |
4.3 参数
| 参数 | 说明 |
|---|---|
spark.sql.streaming.statefulOperator.checkpoint.location | 状态 checkpoint 目录 |
| watermark 设置 | 比真实延迟略大,过小丢数据,过大延迟结果 |
五、状态管理
5.1 有状态操作
| 操作 | 说明 |
|---|---|
| groupBy + 聚合 | 分组聚合(有状态) |
| window 聚合 | 窗口状态 |
| mapGroupsWithState | 自定义状态(强类型) |
| flatMapGroupsWithState | 自定义状态 + 输出控制 |
5.2 状态存储
| 后端 | 说明 |
|---|---|
| HDFS(默认) | 状态存 HDFS,容错 |
| 内存 + WAL | 执行时内存,checkpoint 持久化 |
5.3 状态超时
scala
// 自定义状态设置超时
groupState.setTimeoutDuration("10 minutes")
// 超时触发 onTimeout| 作用 | 说明 |
|---|---|
| 清理无效状态 | 防止状态无限增长 |
| 触发最终结果 | 窗口/会话结束时输出 |
六、输出模式
6.1 三种输出模式
| 模式 | 行为 | 适用 |
|---|---|---|
| Append | 只追加新行 | 无状态/过滤 |
| Update | 更新已输出行 | 聚合、状态 |
| Complete | 全量重写结果 | 全量聚合(结果小) |
6.2 Sink 支持
| Sink | 模式支持 |
|---|---|
| File(Parquet) | Append |
| Kafka | Append/Update |
| Foreach | 全模式(自定义) |
| 内存表 | Append/Update/Complete |
| JDBC(foreach) | 全模式(自定义) |
scala
// Foreach 自定义输出(如写 MySQL)
df.writeStream.foreach(new ForeachWriter[Row] {
override def open(...) = ...
override def process(row: Row) = ...
override def close(...) = ...
}).start()七、容错与 Exactly-Once
| 机制 | 说明 |
|---|---|
| Checkpoint | 状态与 offset 持久化 |
| 源重放 | Kafka 按 offset 重放 |
| 幂等输出 | 结果落 HDFS/Kafka 天然幂等 |
| 端到端 | 源 + 计算 + Sink 三端配合 |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 窗口结果一直不输出 | 水位线设置过大,等待窗口关闭 |
| 迟到数据被丢 | 期望保留则调大 watermark |
| 状态无限增长 | 设置超时/过期策略 |
| Update 模式文件 Sink 不可用 | 文件 Sink 仅 Append |
| 连续模式算子不支持 | 换微批或检查算子支持列表 |