Options
- Mark as New
- Bookmark
- Subscribe
- Mute
- Subscribe to RSS Feed
- Permalink
- Report Inappropriate Content
03-16-2024 10:55 AM
Thanks for your answer @Retired_mod. I have tried what you suggested with the following code in PySpark. I also add some lines to handle scenarios where receiving multiple files/partitions: for example in the first ingestion where I'm reading historical data:
output_location = 'path'
def upsert_to_parquet(batchDF, batchId):
# Extract partition keys (year, month, day) from batchDF
# In case more than one partition is updated iterate over them
partitions = batchDF.select('year', 'month', 'day').distinct().collect()
for partition in partitions:
year, month, day = partition['year'], partition['month'], partition['day']
# Construct the target path based on partition keys
partition_path = f'{output_location}data/year={year}/month={month}/day={day}'
# Overwrite existing data in the target partition
batchDF.filter((col('year') == year) & (col('month') == month) & (col('day') == day)).write.mode('overwrite').parquet(partition_path)
df.writeStream \
.foreachBatch(upsert_to_parquet) \
.option('checkpointLocation', output_location + 'checkpoint') \
.trigger(availableNow = True) \
.start()Apparently it solves the duplication on updates. However there are some additional issues:
- Some extra files appear on my directory on every partition with logs: _committed_xxx, _started_xxx, _SUCCESS files
- For each partition I got two .parquet files: part-0000-xxx.snappy.parquet, part-0001-xxx.snappy.parquet. The first one is empty, with only the schema information (column names). And the second one contains all the data
Ideally only one parquet file should be created. Do you know why this happen?