1. Real-time Mode是什么:先分清两种“实时”
在聊Spark Streaming的实时模式之前,必须先把一个被用滥的词挑明白——“实时”。很多团队跟我聊需求时张口就是“我们要实时数仓”,结果一细问,T+1报表就算实时。真正做流计算的工程师都清楚,所谓的实时,在不同语境下差着十万八千里。
Spark Streaming从早期的DStream模型一路演进到现在的Structured Streaming,最大的变化,就是把“流”从RDD级别的离散批次抽象成了“无界表”(Unbounded Table)。这个抽象的改变是本质性的:你在流上写聚合、写过滤、写Join,跟写批处理SQL几乎一模一样。但要命的是,底层执行方式不一样,延迟表现也不一样。
在Spark 4.0之前,Structured Streaming主要跑在微批模式(Micro-Batch Mode)上。什么叫微批?就是说Spark虽然把数据看成“流”,但实际处理的时候,还是按一小段时间窗口内的数据攒成一个批次来做计算。比如你设置了trigger间隔2秒,那么每2秒,Spark才处理一次这2秒内到达的数据。好处是容错机制可以复用批处理那套成熟方案——WAL(Write-Ahead Log)加数据重放,非常稳。坏处是延迟注定做不到毫秒级,数据从产生到被处理,至少要等一个trigger周期。对于报表、监控等场景这个延迟完全没问题,但比如实时风控、实时竞价这种延迟敏感的业务,2秒就有可能是致命的。
Real-time Mode,也就是Structured Streaming的持续处理模式(Continuous Processing),就是冲着这个短板来的。它把处理模式从“周期性触发”改成“事件驱动”,数据一到马上处理,理论上可以把端到端延迟压到毫秒级。
两类模式的对比我整理了一张表,用的时候直接对着看:
| 对比项 | 微批模式 | 持续处理模式 |
|---|---|---|
| 延迟水平 | 秒级,取决于trigger间隔 | 毫秒级,事件一到即处理 |
| 编程模型 | DataFrame/DataSet API | DataFrame/DataSet API |
| 容错实现 | WAL + 批次重放 | 分布式快照(Chandy-Lamport算法变种) |
| 适用算子 | 几乎所有Structured Streaming算子 | 受限,目前主要支持查询类、过滤类、部分聚合 |
| 生产环境成熟度 | 极高,大规模生产验证充分 | 可用,但复杂状态操作需谨慎评估 |
| 数据一致性 | Exactly-Once保证成熟 | 端到端一致性仍在持续完善 |
我见过不少工程师一上来就追求毫秒级实时,把作业切到持续处理模式,结果发现状态管理、水位线、复杂聚合等一堆问题随之而来,最后又灰溜溜地退回微批。所以先给一个明确的建议:不是所有场景都需要切Real-time Mode。真正适合持续处理的,是链路简单、以过滤转换和简单聚合为主的场景;如果要做窗口聚合、状态管理复杂的业务,微批模式依然是更稳妥的选择。
另外提醒一句,Spark 4.0之后Structured Streaming对持续处理模式的支持已经比早期版本成熟很多,但很多团队还停留在Spark 2.x时代写DStream的老代码上。如果你们还在用JavaStreamingContext那套API,建议尽早规划迁移——DStream虽然灵活,但它和DataFrame API之间的差异越来越大,团队协作成本也在变高。
2. 解锁实时模式的关键能力:事件时间、水印与状态管理
想真正用好Spark Streaming的Real-time Mode,绕不开三个核心概念:事件时间、水印(Watermark)、状态管理。这仨是流计算里的“硬核”,也是和批处理思维差异最大的地方。不把这些搞清楚,就算代码能跑,结果也是错的。
2.1 事件时间与处理时间:流里的两种时间概念
流计算里最常见的一个思维陷阱,是混淆事件时间和处理时间。
- 处理时间(Processing Time):数据被Spark集群真正处理的那个时刻,是机器的系统时间。
- 事件时间(Event Time):数据本身携带的时间戳,比如订单创建时间、日志记录时间、传感器采样时间。
举个例子,一条订单日志在14:30生成,但因为网络抖动、Kafka积压或者其他原因,到15:07才被Spark读到。处理时间是15:07,如果你按处理时间聚合,这条订单会被划进15:00-15:05的窗口,而不是它真正发生的14:30-14:35窗口。对于实时报表、异常检测等场景,这种错配是致命的。
所以,只要业务上关心“这件事什么时候发生的”,就必须用事件时间,并且在代码里显式指定时间字段,让Spark知道哪一列是事件时间。用spark-sql的写法就是在SELECT里对时间列调用一个窗口函数来分组。
这里还要区分一下,“实时”不等于“按到达时间聚合”。很多刚上手Structured Streaming的同学,一边说自己在做实时计算,一边用处理时间去开窗口,这本质上还是批处理思维,只是把批量变小了而已。真正的实时分析,一定要以事件时间为准锚定业务语义。
2.2 水印:怎么容忍乱序与延迟
选定了事件时间后,紧接着的问题就是:数据到达顺序不是按事件时间排好的。有可能14:30的数据比14:35的数据晚到,这就是乱序。流是无穷的,你不可能永远等下去——如果为了等一条迟迟不到的14:30数据而阻塞整个流,那延迟就无底洞了。
水印机制就是用来回答“我到底要等多久”的。它的本质是:声明一个时间阈值,比如“我最多容忍5分钟延迟”。Spark会持续跟踪已经看到的最大事件时间减去这个阈值,称之为水印水位线。水位线是一个时间边界,早于水位线的数据一到,就直接判断为迟到数据,不再参与窗口计算。
代码里设置水印非常简单,核心思想就是事件时间字段定义好后,对event_time列加上一个容忍阈值。水印设置的关键决策是阈值大小:设大了,窗口关闭时间晚,结果准确度高,但延迟大、状态也占内存;设小了,延迟低,但数据就可能被过早丢弃,造成统计偏差。
我的经验是,先观察业务数据的延迟分布,再定水印。比如我们的日志系统99.9%的数据在3分钟内到达,那我就会设4到5分钟的水印,留出缓冲。为了观察这个分布,可以先跑一段时间,在日志里记录到达时间减去事件时间的差值,画个分布图再拍板。别拍脑袋就拿个5分钟,也别信别人“设2分钟都没事”的说法。
2.3 状态管理与容错:实时模式稳不稳,全看这里
持续处理模式能达到毫秒级延迟,关键在于抛弃了“批次”这个中间层。数据是逐条处理的,不需要等攒批。但这也带来了新的问题:如果一个节点在处理某条数据时挂掉了,怎么保证状态一致?微批模式靠WAL重放整个批次,批次边界就是天然的恢复点;持续处理模式没有批次边界,就得靠分布式快照。
分布式快照的基本思路是周期性地把所有任务的状态和执行位置记录到一个checkpoint中。恢复的时候,直接从最近一次快照开始,把没处理完的数据重放一遍。这套机制的原型是Chandy-Lamport算法,Flink用得早,Spark则是在持续处理模式里逐步完善的。
需要注意,checkpoint是一个重操作,不是越频繁越好。频繁checkpoint可以缩短恢复时间,但会显著增加额外的计算和IO开销,压缩了本来想省的时间。生产环境里我会把checkpoint间隔设置在1到2分钟左右,然后配合上游数据源的重放能力来做兜底。好在Structured Streaming天然对接Kafka这类可重放的消息队列,数据源侧的重放能力是够用的,真正要操心的是输出端的幂等性:如果任务重启后重新计算结果,下游能不能接受重复写入?
这里给一个判断清单,属于我用真金白银换来的经验:
- 设置合理的checkpoint间隔,别为了“防止丢数据”没底线地频繁checkpoint。
- 状态管理的内存占用要持续观察,持续处理模式下的状态同样会累积,只增不减的状态最终会把内存吃满。
- 为你的输出端设计幂等机制,这是端到端一致性的最后一公里,Kafka写坏了好办,MySQL炸了可就麻烦了。
- 避免在持续处理模式中混用复杂聚合和Join,这类操作对状态的要求高,目前在持续模式下支持尚不完善,硬上就是给自己找事故。
3. 实操:搭一个基于Structured Streaming的限时统计作业
前面的原理讲清楚了,下面直接上一套可以跑通的实操。我会以“实时订单金额统计”为例,场景设定为:从Kafka消费订单JSON数据,按事件时间的1分钟窗口统计订单总额,结果写入MySQL。这个案例覆盖了Structured Streaming最典型的链路:接入、窗口聚合、输出。
3.1 环境准备工作
我用的是Spark 3.5,Python版本3.10,Kafka 3.4。值得注意的是,从3.2版本以后,spark-sql-kafka-0-10这个包是必须引入的,否则format("kafka")会报Failed to find data source。
依赖方面,如果你是spark-submit提交,加上--packages参数即可;如果在本地IDE调试,把核心的jar包加到工程依赖里就行。另外,MySQL的JDBC驱动也要准备好,在Spark 3.5里,foreachBatch模式下JDBC驱动的类加载问题不常见,但还是要确认驱动包和MySQL版本匹配。
MySQL的表结构也提前建好,别在流作业里建表。表结构类似:窗口开始时间、窗口结束时间、订单总金额、记录更新时间,主键用窗口开始时间和窗口结束时间。
3.2 核心代码骨架与逐段解读
先看完整的Python代码骨架:
from pyspark.sql import SparkSession from pyspark.sql.functions import from_json, col, sum, window, current_timestamp from pyspark.sql.types import StructType, StructField, StringType, DoubleType, TimestampType # 1. 初始化SparkSession spark = SparkSession.builder \ .appName("RealTimeOrderStats") \ .master("local[*]") \ .config("spark.sql.shuffle.partitions", "4") \ .getOrCreate() spark.sparkContext.setLogLevel("WARN") # 2. 定义订单JSON的消息结构 order_schema = StructType([ StructField("order_id", StringType()), StructField("amount", DoubleType()), StructField("event_time", TimestampType()) ]) # 3. 从Kafka消费 raw_df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "localhost:9092") \ .option("subscribe", "order-topic") \ .option("startingOffsets", "latest") \ .load() # 4. 解析JSON并提取事件时间列 order_df = raw_df \ .select(from_json(col("value").cast("string"), order_schema).alias("data")) \ .select("data.*")这段代码里,from_json从Kafka的value字节里解析出结构化字段,然后用select("data.*")把嵌套的JSON字段展开。注意,Kafka的value默认是二进制,需要先.cast("string")再交给from_json,这是一个非常经典的遗漏点。
接着做窗口聚合和水印设置:
# 5. 设置水印并按事件时间开1分钟窗口 windowed_agg = order_df \ .withWatermark("event_time", "2 minutes") \ .groupBy( window(col("event_time"), "1 minute") ) \ .agg( sum("amount").alias("total_amount") ) # 6. 为写入MySQL补充窗口字段 result_df = windowed_agg \ .select( col("window.start").alias("window_start"), col("window.end").alias("window_end"), col("total_amount") )这段是目前为止最关键的部分。withWatermark("event_time", "2 minutes")意味着Spark容忍2分钟内的乱序数据,超过这个范围就被丢弃。groupBy(window(col("event_time"), "1 minute"))按事件时间切1分钟窗口。处理完的效果是:即使数据晚到两三分钟,只要没超过水印,它还是会落到它该落的时间窗里。
写入MySQL的部分用foreachBatch实现:
# 7. 用foreachBatch保证幂等写入MySQL def write_to_mysql(batch_df, batch_id): batch_df.write \ .mode("append") \ .jdbc( url="jdbc:mysql://localhost:3306/realtime", table="order_window_stats", properties={"user": "root", "password": "your_password"} ) # 8. 启动流式计算 query = result_df.writeStream \ .foreachBatch(write_to_mysql) \ .outputMode("append") \ .trigger(processingTime="2 seconds") \ .option("checkpointLocation", "/tmp/spark-checkpoint") \ .start() query.awaitTermination()为什么用foreachBatch而不是writeStream.format("jdbc")直接输出?因为Structured Streaming原生只支持file、kafka、foreach、foreachBatch几种sink。JDBC不属于原生sink,所以用foreachBatch对每个微批做一次批量写,是最稳妥的做法。用foreachBatch还有一个好处是,可以在批次级别对数据做去重、异常值过滤等处理,给下游多一道防护。
还要说明一下输出模式。窗口聚合场景下,outputMode一般用append——新窗口的结果一旦生成就不再修改,直接追加。update模式适合你想实时更新某个key的统计值的场景,complete模式则是每次输出全量聚合结果,在流式场景下要慎用,因为输出数据和压力会随状态增长而无止境膨胀。
关于checkpointLocation,这一项一定要给到可靠的位置。本地测试可以指向本地磁盘,生产环境建议指向分布式文件系统如HDFS,同时保证路径不被级联删除。checkpoint里包含了流的状态、偏移量、元数据等信息,一旦丢失,整个作业的恢复能力就没了。
3.3 关键参数调优参考
调优参数这块我直接给一组可参考的配置值,同时解释每个参数为什么这么设。
| 参数 | 参考值 | 说明 |
|---|---|---|
| spark.sql.shuffle.partitions | 集群核数的2-3倍 | 决定聚合和Join的并行度,太小会热点,太大则小文件多、调度开销大 |
| spark.sql.streaming.schemaInference | true(仅文件源) | 文件源自动推断schema,Kafka源必须显式定义schema |
| spark.sql.streaming.fileSource.schema.forceNullable | false | 防止推断输出全成nullable列,影响下游存储 |
| trigger间隔 | 2-5秒 | 微批模式下延迟与吞吐的折中,生产环境不要小于1秒 |
| spark.sql.adaptive.enabled | true | 开启AQE,动态合并shuffle分区,对倾斜有明显缓解 |
| spark.sql.adaptive.coalescePartitions.enabled | true | AQE的自动合并分区开关 |
对于持续处理模式,trigger的设置方式跟微批不同。持续处理模式用trigger(continuous="2 seconds")这样的写法,但坦白说,持续模式对资源和集群的稳定性要求更高,如果你的集群规模不大、作业拓扑复杂,不要太激进。
4. 常见问题与排查技巧实录
这一节的价值可能比前面代码更值钱。下面每个问题都是我在生产环境里真实踩过、排查过、解决过的,整理成备查清单,希望能帮各位少走弯路。
4.1 作业重启后数据重复或丢失
这是流作业最经典的问题,原因九成在checkpoint。Structured Streaming的故障恢复完全依赖checkpoint中保存的偏移量信息。如果你把checkpoint删了,或者换了checkpoint路径,作业就等于失忆,会从startingOffsets指定的位置重新消费。
排查步骤我建议按下面的顺序来:
- 先确认checkpoint路径是否有效,目录权限是否正确。
- 确认没有多个作业实例共用同一个checkpoint路径。
- 确认
startingOffsets的设置是否符合预期。 - 确认下游写入是否做了主键去重,这是最后一道防线。
我的习惯是,所有下游目标表都设计一个天然的幂等键,比如用window_start + window_end做联合主键,用INSERT ... ON DUPLICATE KEY UPDATE代替纯insert。这样哪怕上游偶尔重复,下游也不会产生脏数据。流作业的可靠性不能只靠处理端,输出端的防御设计同样重要。
4.2 水印设了没效果,过时数据还在
我在排查中见过不少这样的情况:明明设置了withWatermark,但过期数据还是出现在了结果里。原因通常是,你根本没在watermark字段上进行分组。Spark的水印机制只在聚合键包含事件时间列时才会触发过期数据的清理,如果groupBy的字段里只有业务键而没有窗口时间,水印就不会生效。
正确做法是groupBy(window(event_time, "1 minute"), order_id)这样写,让窗口时间参与聚合。水印+窗口必须搭配使用,缺一个都不行。另外,多个聚合流如果做Join,水印的定义要统一,两条流的watermark都要单独设置,并且Join的condition里必须带上event_time的范围条件,否则水印在Join场景下同样不生效。
4.3 持续处理模式下部分算子不支持
切到持续模式,跑着跑着就报Continuous processing does not support ...的错,这是非常正常的。持续处理模式目前支持的算子范围远小于微批模式,凡是需要“攒一批数据才能算”的算子都有一段路要走。比如当前持续模式就没有完整支持任意状态算子。
我的建议是,在做技术选型时就先做一次算子兼容性检查,把要用的算子逐一对照官方文档确认。如果发现有不支持的算子,两条路可选:一是改用微批模式,虽然延迟到秒级,但功能完整;二是把不支持的计算拆出来用其他方式兜底,比如在输出后的下游做二次处理。
务实一点说,大多数业务场景延迟在秒级完全能接受,为了“毫秒级”这个目标去牺牲功能的稳定性,往往得不偿失。我在生产环境里真正跑到持续处理模式的,其实只有两个场景:一个是ETL链路极短的清洗转发,另一个是纯过滤加规则匹配的实时告警。其余场景,微批加合理trigger就够了。
4.4 背压问题:数据积压怎么定位
实时作业最怕的另一个问题是:消费速度赶不上生产速度,Kafka消费组Lag持续上涨。定位它有个清楚的思路:
- 看Kafka的消费组Lag监控,确认积压发生在哪个环节,是数据源读取跟不上,还是处理逻辑本身太慢。
- 看Spark UI的任务耗时分布,是全部任务都慢,还是少数几个任务特别慢——后者大概率是数据倾斜。
- 看执行计划里的shuffle情况,确认是不是某些聚合Key的数据量过大。
处理手段也分优先级:数据倾斜优先用AQE的skewJoin优化,或手动加随机盐;单条处理逻辑慢就检查是否有外部IO调用或序列化开销过大;整个作业吞吐不足则检查分区数和并行度配置。记住一个原则:流作业的吞吐瓶颈,大多数不在CPU而在shuffle和状态IO,优先怀疑这两块。
4.5 状态无限增大导致内存溢出
流作业跑几天之后内存告警,多半是状态只增不减。比如按device_id做累计统计,海量设备ID持续涌进来,状态后端无限膨胀。持续处理模式下,状态保存在内存中,一旦超过Executor内存上限,就会触发了OOM。
应对策略也很明确:
- 检查业务是否真的需要对全量key做累计,如果不是,加时间窗口让过期状态自动清理。
- 检查
watermark设置是否合理,合理的水印能配合窗口释放旧状态。 - 如果某些key确实需要长期保留,那就需要给Executors规划足够的内存,并做好资源隔离。
这块给一个量化经验:实时作业的Executor内存至少要比状态预期峰值再预留30%左右,否则GC频繁带来的延迟抖动会让你想砸电脑。
5. 写在最后的一些实操体会
做流式计算这几年,我最大的感受是:实时不等于快,也不等于准,它是一连串取舍后的结果。你在水印上宽容,数据准了但延迟变高;你设严一点,延迟低了但早到和迟到的数据会有偏差。每个参数背后都是业务诉求在权衡。
如果你刚开始接触Spark Streaming和Real-time Mode,我的建议是,别一上来就追求毫秒级。先把微批模式跑稳,把事件时间、水印、checkpoint、状态管理这些基本功练好,再考虑要不要切持续处理模式。一上来就挑战最高难度的模式,出了问题往往连排查方向都没有。
另外,强烈建议你在本地的Kafka环境里把上文那套代码原样跑一遍,故意用脚本制造乱序数据,观察水印的行为,观察窗口结果的输出时机。这些亲手验证过的东西,比读十篇博客都管用。我在带团队的时候,都会让新人做这两个练习:第一个是乱序数据的窗口统计,第二个是手动kill掉Executor看恢复效果。做完这两个练习,对流计算的理解会上一个台阶。
最后分享一个排障小技巧:Structured Streaming作业先看日志,别急着动参数。WARN级别里隐藏着大量有价值的信息,比如Skipping processing of expired data这类日志,直接暗示水印触发边界出问题了。日志看明白了,问题就解决了一半。