You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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不是用来延迟数据写入的机制,它的核心作用是:

  • 控制流处理中状态的生命周期(自动清理过期状态,避免内存溢出)
  • 在特定操作(窗口聚合、去重)中过滤迟到的重复数据

当天数据被写入的原因

  1. Watermark的阈值逻辑:
    Watermark阈值 = 当前流中观察到的最大startTimestamp - 72小时。对于2023-07-18的当天数据,其事件时间远大于这个阈值,因此不会被Watermark过滤。
  2. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 04:16:31