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.