Spark流处理中Watermark与dropDuplicates的工作机制及数据写入疑问
Spark流Watermark行为分析:为何当天数据仍被写入?
问题背景
设置了72小时的Watermark(基于事件时间startTimestamp),处理窗口为10分钟,流以append模式写入Delta Lake,预期当天(2023-07-18)的最新数据不会写入目标表,但实际查询发现该日期数据已存在。需要明确Watermark有效期内是否会写入数据,以及当前现象的原因。
相关代码片段
流读取与预处理代码
streaming_df = ( self.spark.readStream.format("delta") .option("withEventTimeOrder", "true") .option("skipChangeCommits", "true") .option( "badRecordsPath", f"{self.badrecords_path}/", ) .table(f"{self.catalog_name}.bronze.measurements{self.postfix}") .withWatermark("startTimestamp", "72 hours") .filter((f.col("quality") == "Good") & (f.col("value").isNotNull())) .dropDuplicates(["startTimestamp", "id"]) # ----other transformations )
流写入代码
ws = ( streaming_df.writeStream.queryName("meas_silver_smuk") .foreachBatch( lambda df, epoch_id: process_func( df=df, epoch_id=epoch_id, ) ) .option( "checkpointLocation", f"{checkpoint_path}" ) .trigger(**trigger_options) .start() ) ws.awaitTermination()
问题解答
核心误解纠正
Watermark不是用来延迟数据写入的机制,它的核心作用是:
- 控制流处理中状态的生命周期(自动清理过期状态,避免内存溢出)
- 在特定操作(窗口聚合、去重)中过滤迟到的重复数据
当天数据被写入的原因
- Watermark的阈值逻辑:
Watermark阈值 = 当前流中观察到的最大startTimestamp- 72小时。对于2023-07-18的当天数据,其事件时间远大于这个阈值,因此不会被Watermark过滤。 dropDuplicates结合Watermark的行为:
你使用的dropDuplicates(["startTimestamp", "id"])搭配Watermark,逻辑是:当某个(startTimestamp, id)对应的事件时间早于Watermark阈值时,Spark会判定不会再有该键的迟到数据,于是清理该键的状态并输出去重后的记录。但当天的数据事件时间未超过阈值,只要通过filter条件,就会被实时处理并写入目标表——这是完全符合Spark设计逻辑的。
对文档描述的澄清
文档中提到的Max startTimestamp noticed - Watermark period > End Period of the time window,是针对窗口聚合场景的规则:当窗口的结束时间早于Watermark阈值时,该窗口的聚合结果会被最终确定并输出(因为不会有迟到数据进入该窗口)。但这并不意味着所有数据都要等Watermark到期才写入,普通的过滤、去重操作会实时处理符合条件的数据。
关于dropDuplicatesWithinWatermark
该API是Spark 3.3+引入的,语义上更明确地将去重范围限定在Watermark覆盖的状态内,本质逻辑和你当前的withWatermark + dropDuplicates一致,但可读性更强。
内容的提问来源于stack exchange,提问作者Saugat Mukherjee
相关产品推荐
相关产品推荐

