Pyspark/Databricks按年月分区不新增列且补全空分区方案咨询
Pyspark Databricks Delta表分区问题解法
问题1:消除多余的分区字段
原有代码主动新增了year、month两个列写入表,才会导致这两个字段出现在最终表schema中。submitted_yyyy_mm本身就是标准的年月格式,直接用该字段作为分区键即可,无需拆分,不会新增任何多余字段,写入代码调整如下:
# 无需新增任何额外字段,直接使用现有submitted_yyyy_mm做分区键 df_orders.write \ .partitionBy("submitted_yyyy_mm") \ .mode("overwrite") \ .format("delta") \ .saveAsTable(orders_table)
该方案生成分区的目录格式为submitted_yyyy_mm=2017-01,完全符合年月分区的需求,且支持分区裁剪,查询时筛选年/月效率和两级分区一致。
问题2:生成2017-2019所有36个空分区
Delta Lake支持手动添加无数据的空分区,写入完成后执行以下代码批量添加所有要求的年月分区即可:
# 生成2017-01到2019-12的全量年月序列 all_months = spark.sql(""" SELECT date_format(add_months(to_date('2017-01-01'), seq), 'yyyy-MM') AS month_val FROM sequence(0, months_between(to_date('2019-12-31'), to_date('2017-01-01'))) AS seq """).collect() # 批量添加空分区,已存在的分区会自动跳过 for row in all_months: current_month = row.month_val spark.sql(f"ALTER TABLE {orders_table} ADD IF NOT EXISTS PARTITION (submitted_yyyy_mm = '{current_month}')")
可选兼容方案:保留year/month两级分区
如果业务要求必须使用year/month两级分区结构,又不想这两个字段出现在最终表的查询结果中,可以通过创建视图的方式隐藏分区字段:
from pyspark.sql import functions as F # 生成分区字段写入原始表 df_with_partition = df_orders \ .withColumn("year", F.year(F.col("submitted_yyyy_mm").cast("date"))) \ .withColumn("month", F.month(F.col("submitted_yyyy_mm").cast("date"))) df_with_partition.write \ .partitionBy("year", "month") \ .mode("overwrite") \ .format("delta") \ .saveAsTable(f"{orders_table}_raw") # 创建对外视图,隐藏分区字段 spark.sql(f""" CREATE OR REPLACE VIEW {orders_table} AS SELECT submitted_at, submitted_yyyy_mm, order_id, customer_id, sales_rep_id, shipping_address_attention, shipping_address_address, shipping_address_city, shipping_address_state, shipping_address_zip, ingest_file_name, ingested_at FROM {orders_table}_raw """) # 批量添加两级空分区 all_year_month = spark.sql(""" SELECT year(add_months(to_date('2017-01-01'), seq)) AS year_val, month(add_months(to_date('2017-01-01'), seq)) AS month_val FROM sequence(0, 35) AS seq """).collect() for row in all_year_month: spark.sql(f"ALTER TABLE {orders_table}_raw ADD IF NOT EXISTS PARTITION (year = {row.year_val}, month = {row.month_val})")
内容的提问来源于stack exchange,提问作者Timothy Clotworthy
相关产品推荐
相关产品推荐

