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

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的记录

测试数据说明

  • 表头字段: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 06:12:15