有个朋友前几天问我:“我用 Structured Streaming 写了半天 SQL,也算能跑,但一直没搞清楚它底层每个批次到底在干嘛。批次是怎么触发的?状态到底存到哪了?为什么到处都在说它能做到精确一次?” 我一时不知道从哪讲起,因为这几个问题恰好是 Structured Streaming 和旧版 Spark Streaming、甚至和 Flink 在实现哲学上的分水岭。光看 DataStreamReader、writeStream 这些外层 API,是回答不了这些问题的。
今天这篇是系列第 8 章的附加篇,就专门把这些底层机制拆开来讲。写着之前我特意翻了 Spark 3.x 的源码,也看了线上任务的实际日志,尽量不用“官方文档怎么说”的口吻,而是用“引擎实际怎么跑”的视角。适合正在用 Structured Streaming 做实时任务、想排查线上问题、或者面试前想把微批、状态、容错这些概念讲明白的同学。看完你会对“流处理到底在跑什么”有一个比较立体的认知。
1. 先看设计前提:为什么用“表”来抽象流
1.1 无限表与增量查询
Structured Streaming 最核心的思想就一句话:把流当成一张永远在增长的表格。每来一条数据,就是在表格末尾追加一行。用户写的是 DataFrame 或 SQL,Spark 在逻辑计划里用一个 streaming relation 表示这个数据源,底层对应着 DataSourceV2 的 MicroBatchStream 或 ContinuousStream 接口。
你可能会觉得这只是一个“看着好用”的封装,但实际上它带来的价值比想象中大得多。因为流和批一旦统一成表模型,Catalyst 优化器、Tungsten 执行引擎、WSCG 全代码生成这些批处理积累多年的优化能力,就可以无缝复用到流处理上。你的流任务本质上就是一个无限循环的小型批任务,每个批次拿到的都只是“当前这一小段新增数据”,执行计划却和批处理一模一样。
这里有个容易踩的误区:很多人以为流式 DataFrame 会把全量历史数据攒在内存里。不是的。df.isStreaming返回 true 只代表数据源是流式的,执行层每个批次拿到的是由 offset 划定的那一段数据,不是全表。内存里最多保留状态算子的中间结果,这一点在讲状态管理时还要细说。
1.2 微批模型是刻意选择,不是妥协
Structured Streaming 默认的 Micro-Batch 模式经常被误解成“Spark 做不了实时流”,其实微批是一种工程上的取舍。它把连续到达的数据按触发间隔切成一个个小批次,每个批次本质上就是一个独立的 Spark 作业。好处非常直接:故障恢复只靠记录批次边界和 offset,不需要逐条做分布式快照;调度、内存、shuffle 全都复用 Spark 批处理那套成熟机制;吞吐量能打得很高。
代价也很明显:批次之间的调度开销决定了延迟下限,通常百毫秒到秒级。那为什么不用真正的连续处理?Spark 2.3 确实引入了 Continuous Processing,号称毫秒级延迟,但支持的算子非常有限,只覆盖 select、filter、map 这类无状态操作,做不了聚合和状态管理。生产上绝大多数场景还是微批,如果你需要严格的事件时间窗口和毫秒级延迟,那更适合直接去用 Flink。
2. 一个批次从触发到提交的完整旅程
2.1 触发机制与调度循环
很多人以为 Trigger 就是“每隔多久处理一次”,这个理解不完整。如果完全不指定 Trigger,Spark 的策略是“当前一批处理完成后立刻调度下一批”,只要上游有数据就尽可能快地跑,中间几乎不留空闲时间。使用Trigger.ProcessingTime("10 seconds")时,引擎才会控制批次间隔,而且如果上一批跑了 15 秒,下一批不会并行启动,而是等上一批结束再排队。
这部分的核心执行类是StreamExecution。它内部有一个runActivatedStream循环,每次被触发就执行一批:先获取最新 offset,再调runBatch,等整个批次成功写完 sink 并提交 offset 后,循环才继续。这种单线程的批次调度模型有一个隐藏好处:任务不会因为上游数据堆积而把内存打爆,因为一批没跑完,下一批永远不会进来。
2.2 偏移量管理:WAL 是可靠性的地基
流处理有个经典问题:数据已经处理了,但结果还没写成功,宕机重启以后怎么办。Spark 的答案是先记账再干活。每个 source 在启动时返回一个起始 offset,引擎在每批开始前,先“规划本批要消费的 offset 范围”,然后把这个范围写入 checkpoint 目录下的offsets日志,这个日志就是 WAL。
等 WAL 写成功之后,引擎才真正去 source 拉数据执行计算。这样一来,即使某批拉数失败、计算失败、甚至整个 driver 挂掉,重启后总能从 WAL 找到“本批应该处理哪些数据”,重新拉取再重放一遍。这就是为什么 Structured Streaming 要求必须配置 checkpointLocation——没有 WAL,容错就是空谈。
2.3 执行、输出与提交
本批数据从 source 读进来以后,引擎会结合当前查询的逻辑计划生成一个执行计划,但这个执行计划不是每批重新编译一遍那么重。Catalyst 相当聪明,它会上次生成的物理计划基础上做增量调整,关键路径还是全代码生成的。
计算完成后,结果会写入 sink。写完之后,引擎再往commits目录里写入本批的 batch id,代表这批复数完成。所以 checkpoint 目录里最核心的就是两个文件夹:offsets定义“要处理什么”,commits定义“已经处理完了哪些批次”。恢复时引擎对比两个文件夹,就能知道从哪个 batch 开始重跑。
2.4 端到端精确一次是如何拼出来的
没有任何单一组件能做到端到端精确一次,但组合可以。可重放的 source 加 WAL 保证数据不会丢,版本化 StateStore 保证状态不会算重,幂等或事务性 sink 保证写下游不会重复生效,这几样拼在一起,才构成“端到端精确一次”的完整链路。
| 保障层 | 核心机制 | 解决的问题 |
|---|---|---|
| 数据源 | Kafka 等可重放 Source | 同一段数据能按 offset 重新消费 |
| 批边界 | WAL(offsets 日志) | 宕机后知道该从哪一批继续 |
| 状态存储 | 版本化 StateStore | 重复计算时状态覆盖而不是累加 |
| 结果写出 | 幂等 Sink / 事务 Sink | 重复写入不会产生重复效果 |
这也是官方文档一直强调的那句话:只有在 sink 具备幂等性或事务性的前提下,才能保证端到端精确一次。sink 如果是乱写的普通文件,那最多只能是至少一次。
3. 有状态计算的“记忆”存在哪里
3.1 StateStore 不是内存里的 HashMap 那么简单
有状态算子,比如聚合、去重、flatMapGroupsWithState,微批执行时都需要维护每个分组或每个 key 的中间结果。这个状态由StateStore组件管理,生产环境默认实现是HDFSBackedStateStoreProvider,看得出来,它背后还是分布式文件系统。
每个批次对状态的修改会以 delta 文件的形式写到 checkpoint 目录的state文件夹里,查询时引擎会读取“上一批版本的快照加 delta”合并到内存,批处理完成后再写出新的 delta。为了防止 delta 无限增长,引擎会在积累到一定版本数之后,把增量合并成一次 snapshot 快照。整个流程类比一下就是:你每天写日记(delta),每周把日记整理成年鉴(snapshot),恢复的时候先读年鉴再补几篇日记,比从头读所有日记快得多。
3.2 Watermark:事件时间与状态清理的标尺
为什么官方要求聚合操作必须配 watermark?因为不清理状态,状态会无限膨胀。watermark 的计算方式是“当前观察到的最大事件时间减去延迟容忍度”。举个例子,withWatermark("eventTime", "10 minutes"),假设当前批次里最大事件时间是 12:30,那么 watermark 就是 12:20。事件时间晚于 watermark 的迟到数据,仍然可以参与计算;一旦小于 watermark,就会被丢弃,同时引擎可以安全地清理过期窗口和分组状态。
这里有三个容易混淆的点,我踩过坑才搞清楚。第一,watermark 只在有聚合或状态算子时才有意义,纯 filter 和 select 不需要它。第二,watermark 依赖的是事件时间,也就是数据里的业务字段,不是系统时间。第三,watermark 是单调递增的,它只会往后推,不会往回退。
3.3 分组状态的存储与超时
mapGroupsWithState和flatMapGroupsWithState是比聚合更灵活的状态接口,每个分组 key 可以维护任意自定义对象。状态更新和超时判断都在用户回调函数里完成,所以引擎不知道你的状态结构是什么,它只负责帮你存和取。
超时模式有两种。Processing time timeout 按系统时间判断,适合“订单 15 分钟未支付自动关闭”这类业务,不依赖数据流本身。Event time timeout 依赖 watermark 推进,适合会话窗口这类基于事件时间的逻辑。注意,使用 event time timeout 必须先设置 watermark,否则状态永远不会超时,这个真有人踩过。
3.4 状态膨胀监控与关键参数
每个批次的StreamingQueryProgress里有一个stateOperators数组,里面有numRowsTotal代表状态总行数,numRowsUpdated代表本批更新行数,commitTimeMs代表状态提交耗时。生产上如果发现numRowsTotal持续增长,优先检查是不是没设 watermark、窗口保留时间太长、或者业务 key 太散。
相关参数也要知道:spark.sql.streaming.stateStore.minTotalDirs影响 snapshot 文件在目录间的分布,spark.sql.streaming.stateStore.maxVersions控制 delta 文件保留的版本数。一般不用动,但状态量级上来以后,快照合并频率和文件数会成为潜在瓶颈。
4. 输出模式与结果表机制
4.1 Append、Update、Complete 怎么选
结果表这个概念必须建立起来:用户查询的“结果”在逻辑上也是一张不断演变的表,每个批次执行完,这张表都会发生变化。三种输出模式回答的问题是“本批结束后,把结果表中哪些行发送给 sink”。
Append 模式只输出本批新增的行,适合结构上只会追加不会变更的结果。注意,聚合查询如果选 Append,必须配 watermark,引擎才能判断“哪些窗口已经不可能再更新”,否则会直接报错。Update 模式把本批新增和更新的行一起输出,适合聚合、实时看板、状态变更这类“只关心变化”的场景。Complete 模式每次输出整个结果表,只适合聚合查询,而且要确保结果表本身不大。
4.2 Sink 侧的写入语义
FileSink、KafkaSink、ForeachWriter、ForeachBatch 的保证级别各不相同。FileSink 通过先写临时文件再原子 rename 的方式,保证单个文件不会半截可见。KafkaSink 其实不是事务性的,写 Kafka 本身是至少一次,需要下游幂等消费配合去重。ForeachWriter 提供 open、process、close 三个回调,你可以基于批次和分区自己做幂等提交。ForeachBatch 是给那些需要整批处理的外部系统用的,直接把 DataFrame 交到你手里,由你决定怎么写。
这里给一个选型建议:能基于主键 upsert 的下游,比如 MySQL、Doris、ES,优先用foreachBatch配合批量写入;下游只支持追加的,比如 Kafka、HDFS,用原生 sink 或 writer 更省心。
4.3 从进度日志里看每一批到底干了什么
打开 INFO 级别日志,能看到类似这样的输出:
{ "id" : "63f4c1a2-...", "runId" : "8a0b2c1e-...", "batchId" : 78, "numInputRows" : 1200, "inputRowsPerSecond" : 300.0, "processedRowsPerSecond" : 250.0, "durationMs" : { "addBatch" : 120, "getBatch" : 80, "queryPlanning" : 20, "walCommit" : 15 }, "stateOperators" : [ { "numRowsTotal" : 50000, "numRowsUpdated" : 300 } ] }这段 JSON 里最有诊断价值的是durationMs。getBatch大说明 source 拉数或反序列化慢,addBatch大说明计算、shuffle、状态读写或 sink 写入慢,queryPlanning大说明计划生成阶段有问题但很少见。我排查线上任务时,基本靠连续观察几个批次的这几个指标,就能把瓶颈定位到具体环节。
5. 实操:做一个“订单累计金额”实时统计并验证机制
5.1 需求与方案选型
假设 Kafka 里有订单事件,字段包括 userId、orderId、amount、eventTime,现在要统计每个用户累计下单金额,并且金额变化时要实时更新到 MySQL。这个场景不需要开窗口,因为统计的是全量累计值,所以代码里不写 window,也不强制配 watermark。但要注意,无窗口持续聚合的状态量等于用户数,如果用户量特别大,状态膨胀是必然的,需要定期清理不活跃用户或者分层存储。
输出模式选 Update 而不是 Append。因为每组数据更新时,下游要拿到“这个用户的最新累计金额”。如果选 Append,聚合结果在 Spark 看来是“旧结果不变,新结果追加”,这不符合累计更新的语义,实际跑的时候也会报错。
5.2 关键代码与参数说明
完整伪代码如下:
val df = spark.readStream .format("kafka") .option("kafka.bootstrap.servers", "localhost:9092") .option("subscribe", "order_topic") .option("startingOffsets", "latest") .load() val orderStream = df .selectExpr("CAST(value AS STRING) as json") .select(from_json(col("json"), orderSchema).as("data")) .select("data.userId", "data.amount", "data.eventTime") val aggResult = orderStream .groupBy("userId") .agg(sum("amount").as("totalAmount")) val query = aggResult .writeStream .outputMode(OutputMode.Update()) .trigger(Trigger.ProcessingTime("5 seconds")) .option("checkpointLocation", "hdfs:///checkpoint/order-total") .foreachBatch { (batchDf, batchId) => batchDf.write .format("jdbc") .option("url", "jdbc:mysql://...") .option("dbtable", "user_order_total") .option("user", "...") .option("password", "...") .mode(SaveMode.Append) .save() } .start() query.awaitTermination()foreachBatch 里我用 MySQL 的ON DUPLICATE KEY UPDATE做 upsert,这样同一批数据即使因为故障被重放,最终也只保留正确的累计值,不会累加两次。trigger 设 5 秒是为了观察批次行为,生产上一般根据上游延迟和下游吞吐来调整,没有固定标准。
5.3 通过 UI 和日志观察微批行为
任务启动后,Spark UI 里有一个 Structured Streaming 标签页,可以看到每个 query 的批次列表、输入行数、两个吞吐指标、调度延迟。另外日志里每批结束都会打印一行Streaming query made progress,包含完整的 JSON,直接 grepbatchId就能看到连续批次的执行指标。
比如说你看到一个批次getBatch占了 500ms,另一个批次addBatch占了 2s,那重点就应该放在状态读写和 MySQL 写入上,而不是去调 Kafka 的拉取参数。UI 和日志配合着看,能少做很多无用功。
5.4 验证精确一次语义
任务跑起来后,可以故意做一次故障演练:先消费一批数据,等 MySQL 里有部分累计值,然后手动 kill 掉 driver,再用同一个 checkpoint 目录重启。重启后 Spark 会从 WAL 里找到未提交的批次,重新拉取重放,最终 MySQL 里的累计值不会因为重复处理而变大。
这个验证结果很大程度上取决于 foreachBatch 里的写入逻辑。如果 MySQL 表是把 userId 设为主键并进行 upsert,重复跑就没有副作用;如果每次只是 insert,那么至少一次语义下就会出现重复记录。所以“精确一次”从来不是流引擎单方面保证的,sink 侧必须配合。
6. 常见问题与排查技巧实录
6.1 状态膨胀导致任务越来越慢
症状是处理延迟持续上升,日志里stateOperators.numRowsTotal越涨越高。原因基本逃不出这几类:没设 watermark、key 粒度过细、超时逻辑没生效、watermark 没有随着事件时间推进。排查第一步先看numRowsTotal有没有下降趋势,第二步分析状态 key 的分布,第三步检查事件时间字段取值是否正常。如果业务上允许,可以考虑用两层聚合加盐,把单 key 状态拆散到多个子 key。
记住一个原则:状态量的增长模型必须提前设计好,不要等任务 OOM 了才去翻日志,那时候已经晚了。
6.2 输入数据倾斜导致状态写热点
按用户分组时,如果某些大用户的数据量远高于其他人,那这些 key 的状态更新会集中在少数 executor 上,表现为部分节点 CPU 打满、shuffle 数据量不均衡。常用策略是加盐:在 key 后拼接随机后缀拆成多个子 key,做第一层聚合,下游再按真实 key 做第二层合并。代价是状态量会变大,但每个节点的写入压力被摊平了,适合热点非常明显的场景。
6.3 任务启动后没有任何输出
先确认numInputRows是不是 0。如果是,说明 source 就没拉到数据,问题在 Kafka offset 或消费组,而不是计算逻辑。如果不是 0 但结果迟迟不出现,大概率是输出模式或窗口状态的问题:聚合加 Append 没配 watermark 会直接报错,配了 watermark 也要等窗口关闭才输出第一行。还有一种情况容易被忽略:JDBC 连接池配置太小,批量写 MySQL 时任务在刷日志重试,表面看就是“卡住了”。
6.4 checkpoint 恢复失败
这是生产环境最头痛的问题。Structured Streaming 的 checkpoint 兼容性并不像想象中那么宽松,任何改变查询计划的操作,比如调整 groupBy 的 key、修改 watermark delay、增加中间字段,都有可能让反序列化 StateStore 时抛异常。经验做法是:生产代码不要轻易改状态定义;非改不可时,评估状态丢失的代价,宁可新建 checkpoint 从源头重放一部分数据,也比在一个不兼容的旧 checkpoint 上反复试错强。
6.5 吞吐量上不去
下面这张表是我常用的调优检查顺序,按优先级从高到低排列:
| 检查项 | 操作 |
|---|---|
| Kafka 分区数 | 确保分区数不少于执行端并行度 |
| source 并行度 | 使用minPartitions提高读分区数 |
| shuffle 分区数 | 调整spark.sql.shuffle.partitions |
| 单批大小 | 用maxOffsetsPerTrigger限制一批的数据量 |
| 批次间隔 | 适当加大 trigger 间隔,降低状态提交频率 |
调参之前先看processedRowsPerSecond和durationMs的分布,不要盲目加资源。很多时候瓶颈根本不在计算,而在状态提交或者下游写入的吞吐上限,加再多 executor 也没用。
这个内容后续还可以扩展的方向有两个,一个是 Continuous Processing 模式下的执行机制,另一个是状态存储从 HDFS 迁移到 RocksDB 的实现细节。不过那两个场景比较偏,大部分做实时数据开发的同学可能一辈子都用不上。还是先把微批的这套执行链路吃透,线上问题能定位、面试能讲清楚,就已经非常够用了。