Spark Structured Streaming中PySpark与Scala水印行为差异问题
PySpark Structured Streaming Watermark 异常过滤问题记录
问题背景
使用Python开发搭载Watermark机制的Spark Structured Streaming应用,采用console sink开展测试时,发现PySpark与Scala API存在不符合预期的行为差异:
- 测试场景将派生列作为窗口与Watermark依赖的event-time列
- 测试使用两份数据文件:
- 文件1包含时间戳为
2022-06-05 22:20:16的记录 - 文件2包含时间戳为
2022-06-05 23:35:16的记录
- 文件1包含时间戳为
测试数据说明
- 表头字段:
originalTimestamp,balance - 文件1内容:
2022-06-05 22:20:16,5 - 文件2内容:
2022-06-05 23:35:16,7
核心处理逻辑
streamingDF = spark.readStream.schema( customSchema).csv("csv_path").withColumn("ts", to_timestamp("originalTimestamp")) interDF = streamingDF.withWatermark("ts", "1 minute").groupBy( window(streamingDF.ts, "2 minutes")) outDF = interDF.agg(sum(col("balance")).alias("threshold")).filter(col("threshold") > 5)
流输出实现代码
PySpark版本
query = ( outDF.writeStream.format("console").outputMode("update"). option("checkpointLocation", "checkpoint_folder").start()) query.awaitTermination()
Scala版本
outDF.writeStream .format("console") .option("checkpointLocation", "checkpoint_folder") .outputMode("update") .start() .awaitTermination()
测试执行步骤
- 分别独立启动PySpark与Scala版本的流处理作业
- 先后两次将文件1放入作业监控的CSV路径(使聚合结果满足threshold>5的过滤条件),两个版本作业均正常输出计算结果
- 放入文件2,此时两个作业checkpoint路径中记录的watermarkMs均为
1654468456000 - 再次放入文件1,此时系统计算的水印位置为
23:34:16,传入记录的事件时间为22:20:16,按照Watermark机制该迟到数据理应被过滤丢弃
实际运行结果
- Scala版本运行符合预期,未处理该过期迟到数据
- PySpark版本在后续批次中错误拾取了该过期数据并输出了聚合结果,不符合Watermark机制的设计预期
已完成排查操作
- 将PySpark版本从3.2.1升级至最新的3.3.0版本,问题仍复现
- 校验watermarkMs取值,确认其为当前已处理批次中的最大事件时间,取值逻辑正确
- 暂未定位到该问题的根因,征集同类问题的排查思路与解决方案
运行输出对比
PySpark运行输出
------------------------------------------- Batch: 0 ------------------------------------------- +--------------------+---------+ | window|threshold| +--------------------+---------+ |{2022-06-05 22:20...| 10.0| +--------------------+---------+ ------------------------------------------- Batch: 1 ------------------------------------------- +------+---------+ |window|threshold| +------+---------+ +------+---------+ ------------------------------------------- Batch: 2 ------------------------------------------- +--------------------+---------+ | window|threshold| +--------------------+---------+ |{2022-06-05 23:34...| 7.0| +--------------------+---------+ ------------------------------------------- Batch: 3 ------------------------------------------- +------+---------+ |window|threshold| +------+---------+ +------+---------+ ------------------------------------------- Batch: 4 `The batch that is not right!!` ------------------------------------------- +--------------------+---------+ | window|threshold| +--------------------+---------+ |{2022-06-05 22:20...| 15.0| +--------------------+---------+
Scala运行输出
------------------------------------------- Batch: 0 ------------------------------------------- +--------------------+---------+ | window|threshold| +--------------------+---------+ |[2022-06-05 23:34...| 7.0| |[2022-06-05 22:20...| 10.0| +--------------------+---------+ ------------------------------------------- Batch: 1 ------------------------------------------- +------+---------+ |window|threshold| +------+---------+ +------+---------+ ------------------------------------------- Batch: 2 ------------------------------------------- +------+---------+ |window|threshold| +------+---------+ +------+---------+ ------------------------------------------- Batch: 3 ------------------------------------------- +------+---------+ |window|threshold| +------+---------+ +------+---------+ ------------------------------------------- Batch: 4 ------------------------------------------- +------+---------+ |window|threshold| +------+---------+ +------+---------+
内容的提问来源于stack exchange,提问作者Moh89
相关产品推荐
相关产品推荐

