Spark写入Delta表时partitionBy分区数量不符问题咨询
核心原因:你混淆了两个完全不同的分区概念
df.rdd.getNumPartitions()返回的是Spark计算引擎层面的RDD并行分区数,代表这份DataFrame在内存中被切分为15个分片做并行计算,和你写入Delta表时按date字段拆分的表分区没有直接关联。- 你预期的32个按
date拆分的是Delta表的存储分区(目录分区),统计该值不要用RDD分区数判断,可执行以下命令查看:
-- 查看Delta表的所有存储分区 SHOW PARTITIONS `table`; -- 或者查看表详情确认分区字段和分区数量 DESCRIBE DETAIL `table`;
Python环境下也可以通过Delta API查询:
from delta.tables import DeltaTable delta_table = DeltaTable.forName(spark, "table") # 查看所有date唯一值,对应存储分区数量 print(delta_table.toDF().select("date").distinct().count())
为什么会出现该差异
Spark的RDD分区数由引擎自身的并行规则决定,和业务分区字段的唯一值数量无关:
- 读取数据源时的默认分片规则,比如大量小文件会被自动合并减少分区数
- 若DataFrame执行过shuffle操作,分区数默认由
spark.sql.shuffle.partitions参数控制,开启AQE(自适应查询执行)后还会自动合并小分区,数据量较小时分区数会远低于默认的200 - 写入Delta表时,1个RDD分区可以包含多个
date值的数据,Spark会自动将同一date的数据写入对应的存储分区目录,不会因为RDD分区数少就合并存储分区,正常情况下你最终还是会得到32个存储分区。
如果实际查询到存储分区确实只有15个,可排查以下问题
- 检查
date字段是否存在脏数据,比如空值、格式不一致(如2023/01和2023-01被判定为不同值,或部分值解析失败为同一默认值) - 确认写入前DataFrame的
date唯一值数量是否真的为32,执行df.select("date").distinct().count()验证 - 检查写入前是否有过滤逻辑误删了部分
date对应的数据
内容的提问来源于stack exchange,提问作者SVG
相关产品推荐
相关产品推荐

