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

Spark Structured Streaming watermark失效:旧数据仍输出问题排查

Spark Structured Streaming Watermark机制未生效问题排查

问题描述

编写Spark Structured Streaming代码实现watermark机制处理延迟数据,发现watermark被忽略:即使旧数据的事件时间戳超出(最大事件时间戳 - watermark)范围,仍会被处理并输出。例如最后一条输入数据的ts为2022-03-14T16:16:00.000-07:00,延迟2天仍被输出。

代码示例

schema = StructType([
            StructField("temparature", LongType(), False),
            StructField("ts", TimestampType(), False),
            StructField("insert_ts", TimestampType(), False)
        ])


streamingDataFrame = spark \
                .readStream \
                .format("kafka") \
                .option("kafka.bootstrap.servers", kafkaBrokers) \
                .option("group.id", 'watermark-grp') \
                .option("subscribe", topic) \
                .option("failOnDataLoss", "false") \
                .option("includeHeaders", "true") \
                .option("startingOffsets", "latest") \
                .load() \
                .select(from_json(col("value").cast("string"), schema=schema).alias("parsed_value"))

resultC = streamingDataFrame.select( col("parsed_value.ts").alias("timestamp") \
                   , col("parsed_value.temparature").alias("temparature"), col("parsed_value.insert_ts").alias("insert_ts"))


resultM = resultC. \
    withWatermark("timestamp", "10 minutes"). \
    groupBy(window(resultC.timestamp, "10 minutes", "5 minutes")). \
    agg({'temparature':'sum'})

resultMF = resultM. \
            select(col("window.start").alias("startOfWindowFrame"),col("window.end").alias("endOfWindowFrame") \
                          , col("sum(temparature)").alias("Sum_Temperature"))

result = resultMF. \
                     writeStream. \
                     outputMode('update'). \
                     option("numRows", 1000). \
                     option("truncate", "false"). \
                     format('console'). \
                     option('checkpointLocation', checkpoint_path). \
                     queryName("sum_temparature"). \
                     start()

result.awaitTermination()

输入Kafka数据

+---------------------------------------------------------------------------------------------------+----+
|value                                                                                              |key |
+---------------------------------------------------------------------------------------------------+----+
|{"temparature":7,"insert_ts":"2023-03-16T15:32:35.160-07:00","ts":"2022-03-16T16:12:00.000-07:00"} |null|
|{"temparature":15,"insert_ts":"2023-03-16T15:33:24.933-07:00","ts":"2022-03-16T16:12:00.000-07:00"}|null|
|{"temparature":11,"insert_ts":"2023-03-16T15:37:36.844-07:00","ts":"2022-03-15T16:12:00.000-07:00"}|null|
|{"temparature":8,"insert_ts":"2023-03-16T15:41:33.312-07:00","ts":"2022-03-16T10:12:00.000-07:00"} |null|
|{"temparature":14,"insert_ts":"2023-03-16T15:42:27.627-07:00","ts":"2022-03-16T10:10:00.000-07:00"}|null|
|{"temparature":6,"insert_ts":"2023-03-16T15:44:44.508-07:00","ts":"2022-03-16T11:16:00.000-07:00"} |null|
|{"temparature":19,"insert_ts":"2023-03-16T15:46:15.486-07:00","ts":"2022-03-16T11:16:00.000-07:00"}|null|
|{"temparature":3,"insert_ts":"2023-03-16T16:10:15.676-07:00","ts":"2022-03-16T16:16:00.000-07:00"} |null|
|{"temparature":13,"insert_ts":"2023-03-16T16:11:52.194-07:00","ts":"2022-03-14T16:16:00.000-07:00"}|null|
+---------------------------------------------------------------------------------------------------+----+

输出结果

-------------------------------------------
Batch: 14
-------------------------------------------
+-------------------+-------------------+---------------+
|startOfWindowFrame |endOfWindowFrame   |Sum_Temperature|
+-------------------+-------------------+---------------+
|2022-03-16 16:15:00|2022-03-16 16:25:00|3              |
|2022-03-16 16:10:00|2022-03-16 16:20:00|3              |
+-------------------+-------------------+---------------+

-------------------------------------------
Batch: 15
-------------------------------------------
+-------------------+-------------------+---------------+
|startOfWindowFrame |endOfWindowFrame   |Sum_Temperature|
+-------------------+-------------------+---------------+
|2022-03-14 16:15:00|2022-03-14 16:25:00|13             |
|2022-03-14 16:10:00|2022-03-14 16:20:00|13             |
+-------------------+-------------------+---------------+

问题根源与修正

核心错误

代码中groupBy操作使用了原始DataFrame resultC.timestamp,而非经过withWatermark处理后的列。Spark无法将watermark规则与窗口聚合关联,导致watermark机制完全失效。

修正后的代码

将groupBy中的resultC.timestamp替换为col("timestamp")(或直接写"timestamp"),确保使用withWatermark后的列:

resultM = resultC. \
    withWatermark("timestamp", "10 minutes"). \
    groupBy(window(col("timestamp"), "10 minutes", "5 minutes")). \
    agg({'temparature':'sum'})

原理说明

Watermark的生效依赖两个条件:

  1. withWatermark必须在聚合操作之前调用
  2. 聚合操作使用的事件时间列必须与withWatermark指定的列一致

满足条件后,Spark会自动跟踪最大事件时间,当窗口的结束时间小于max(event_time) - watermark时,会清理该窗口的状态,后续到达的延迟数据将被丢弃,不会被处理。

补充问题解答:如何确保超过watermark的延迟数据不被处理

正确配置watermark后,Spark会自动过滤并丢弃超出watermark范围的延迟数据,但需注意以下几点:

  • 严格遵循操作顺序:先withWatermark,再基于同一事件时间列做窗口聚合
  • 若需在聚合前主动过滤延迟数据,可在withWatermark后添加过滤逻辑(依赖系统时间,不如watermark精准):
    .filter(col("timestamp") >= current_timestamp() - expr("INTERVAL 10 MINUTES"))
    
  • 确保checkpoint目录正确配置,状态管理正常,否则过期窗口状态无法被清理,延迟数据仍可能被处理

内容的提问来源于stack exchange,提问作者Karan Alang

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 16:27:59