spark streaming watermark and explode error, file: .../state/0/0/1.delta does not exist

Viewed 73

I'm trying to write a dataframe to S3 in append mode, I grouped the data by the id and the a window for the markdown, and then I applied explode to one of the fields, using the following code:

s3Session = df
             .groupBy(col("id"), window(col("timestamp"), "1 day"))
             .agg(
                   collect_list(struct($"page", $"pageIndex", $"event")).as("visitedPages")
             )
             .select($"id", $"window", explode($"visitedPages"))
             .select($"id", $"window", $"col.page", $"col.pageIndex")

then I tried to write the result into S3 by:

 s3Session
  .writeStream
  .format("parquet")
  .outputMode("append")
  .option("checkpointLocation", "/tmp/s3Checkpoint")
  .option("path", S3_uri)
  .start()
  .awaitTermination()

This is the error I'm receiving:

22/01/08 10:05:30 WARN HDFSBackedStateStoreProvider: The state for version 9 doesn't exist in loadedMaps. Reading snapshot file and delta files if needed...Note that this is normal for the first batch of starting query.
22/01/08 10:05:30 WARN HDFSBackedStateStoreProvider: The state for version 9 doesn't exist in loadedMaps. Reading snapshot file and delta files if needed...Note that this is normal for the first batch of starting query.
22/01/08 10:05:30 ERROR Executor: Exception in task 0.0 in stage 1.0 (TID 1)
java.lang.IllegalStateException: Error reading delta file file:/tmp/s3Checkpoint/state/0/0/1.delta of HDFSStateStoreProvider[id = (op=0,part=0),dir = file:/tmp/s3Checkpoint/state/0/0]: file:/tmp/s3Checkpoint/state/0/0/1.delta does not exist
....

Keep in mind the code was working properly and writing to S3 before reaching this step and using explode.

0 Answers
Related