I want to generate parquet files every hour with all the information received during that hour for further processing using Spark NLP.
I have streaming data coming from Kafka and when I writeStream I set the trigger processing time to one hour, it just generates many small parquet files. I've read that setting coalesce to 1 plus the trigger would generate a big parquet file but it still gives me many small parquet files.
I've also read that one can also set the minimum number of rows too, but the amount of rows I receive changes from hour to hour.
This is how I'm writing my stream:
df.writeStream \
.format("parquet") \
.option("checkpointLocation", "s3a://datalake-twitter-app/spark_checkpoints/") \
.option("path", "s3a://datalake-twitter-app/raw_datalake/") \
.trigger(processingTime='60 minutes') \
.start()
Any idea on how I can write one hour long parquet files with all the information received during the given hour using Spark Structured Streaming? Maybe I should use something different?
I thought that it could also be that is better to focus on the reading files from the S3 bucket on the other Spark process, and read the files within the previous hour.