I am quite new to Databricks Delta table and I am trying to figure out how the partition would work performance wise. I understand the why on partitioning, yet not exactly the how.
As an example, consider the following (completely arbitrary) data
from pyspark.sql import *
import datetime
City = Row("Name", "CovidInfections", "Date")
city1 = City("Amsterdam", 312, datetime.date(2021, 9, 8))
city2 = City("Amsterdam", 289, datetime.date(2021, 9, 9))
city3 = City("London", 412, datetime.date(2021, 9, 8))
city4 = City("London", 403, datetime.date(2021, 9, 9))
city5 = City("Paris", 201, datetime.date(2021, 9, 8))
city6 = City("Paris", 188, datetime.date(2021, 9, 9))
city7 = City("Berlin", 48, datetime.date(2021, 9, 8))
city8 = City("Berlin", 112, datetime.date(2021, 9, 9))
# Create the dataframe
cities = spark.createDataFrame([city1, city2, city3, city4, city5, city6, city7, city8])
I can write this data to a delta table using a partitioning on Name and Date:
cities.write.partitionBy("Name", "Date").format("delta").mode("overwrite").option("overwriteSchema", True).saveAsTable("covid_cities")
And I see in my data storage separate folders appearing for both Name and Date. However, if I would change data in only one of these partitions, let's say
from pyspark.sql.functions import when
cities_updated = cities.withColumn("CovidInfections", when(cities["Name"]=="Berlin",230).otherwise(cities["CovidInfections"]))
and I would save this updated data
cities_updated.write.partitionBy("Name", "Date").format("delta").mode("overwrite").option("overwriteSchema", True).saveAsTable("covid_cities")
I would expect that only the observations related to Berlin would be affected. However, in my file storage, I can see that a new version of the data is created for every city.
Could anyone elaborate on why this is happening? Am I in fact partitioning right? Or should I save my data in some different way?