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的生效依赖两个条件:
withWatermark必须在聚合操作之前调用- 聚合操作使用的事件时间列必须与
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
相关产品推荐
相关产品推荐

