Spark Structured Streaming - Watermarks behaviour between PySpark and Scala

Viewed 61

I am trying to deploy a structured streaming app in python using watermark. I used console sink to test this, but found weird behaviour between pyspark and scala. I am using a derived column as event-time column on which window and watermark relies. Using two files one with value 2022-06-05 22:20:16 and other with 2022-06-05 23:35:16.

Header: originalTimestamp,balance

File1:

2022-06-05 22:20:16,5

File2:

2022-06-05 23:35:16,7

Logic:

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 (write):

query = (
    outDF.writeStream.format("console").outputMode("update").
    option("checkpointLocation", "checkpoint_folder").start())

query.awaitTermination()

Scala (write):

 outDF.writeStream
    .format("console")
    .option("checkpointLocation", "checkpoint_folder")
    .outputMode("update")
    .start()
    .awaitTermination()

Steps performed:

  1. Run the pyspark and scala separately.
  2. Dropped the first file to location twice (to satisfy filter condition) and see output on both scala and python
  3. Dropped second file (By this time the watermarkMs on both the offsets location is 1654468456000)
  4. Again dropped first file, at this point the watermark is at 23:34:16 but the incoming record has 22:20:16 and it should be neglected. right?

However Scala does the job perfectly - it didn't picked it up! But Pyspark job picked this in next batch and printed the result. Its kind of weird !

Troubleshooting steps performed:

  1. Update pyspark version to 3.3.0 (latest) from 3.2.1
  2. validated the watermarkMs and it is the highest value in latest batch.

Not sure why i am seeing this, anyone has seen this behaviour before? any help would be appreciated!

Output of 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|
+--------------------+---------+

Output of 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|
+------+---------+
+------+---------+
0 Answers
Related