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:
- Run the pyspark and scala separately.
- Dropped the first file to location twice (to satisfy filter condition) and see output on both scala and python
- Dropped second file (By this time the watermarkMs on both the offsets location is
1654468456000) - Again dropped first file, at this point the watermark is at
23:34:16but the incoming record has22:20:16and 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:
- Update pyspark version to 3.3.0 (latest) from 3.2.1
- 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|
+------+---------+
+------+---------+