PySpark按customerid分区写入且不删除历史客户分区的最优方案
最优实现方案
方案1:开启动态分区覆盖(性能最优,无额外开销)
这是Spark原生支持的特性,完美解决overwrite模式删除未更新客户历史分区的问题,全程分布式运行,不需要把客户列表拉取到Driver端,性能和常规批量写入完全一致。
操作步骤
- 先调整Spark会话配置,开启分区动态覆盖模式:
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "DYNAMIC")
该配置的作用是:写入时仅覆盖本次数据集内存在的分区,其余历史分区完全保留,不会被删除。
2. 正常按Customerid分区写入即可:
dfcustomer.write \ .mode("overwrite") \ .partitionBy("Customerid") \ .format("parquet") \ .save("你的目标存储路径")
效果说明
- 写入后目录结构为
目标路径/Customerid=C1/xxx.parquet、目标路径/Customerid=C2/xxx.parquet,天然按客户隔离,每个客户的文件都在对应独立目录下 - 本次数据集没有的客户,对应的历史分区目录和文件完全不会被修改
- 全程仅触发一次Spark作业,无额外的重复扫描、过滤开销,性能远高于循环过滤的写法
方案2:自定义输出文件名包含Customerid(无分区文件夹场景)
如果你不需要Customerid=xxx的分层目录,而是希望所有文件都放在同一级目录,且文件名直接带客户ID,可以按以下方式实现:
# 按Customerid重分区,保证同一个客户的所有数据都在同一个分区文件内 df_repartition = dfcustomer.repartition("Customerid") # 写入时指定文件名规则(Spark 3.0+支持) df_repartition.write \ .mode("overwrite") \ .option("partitionOverwriteMode", "dynamic") \ .option("fileNamePattern", "customer_${Customerid}_${spark_partition_id}.parquet") \ .format("parquet") \ .save("你的目标存储路径")
原有方案性能差的原因
你之前的循环过滤写法需要先把全量客户ID拉取到Driver端,然后对每个客户单独触发一次全表扫描过滤作业,相当于全表扫描次数等于客户数量,数据量较大时会产生大量重复IO和调度开销,性能自然极差。
内容的提问来源于stack exchange,提问作者PysparkSan
相关产品推荐
相关产品推荐

