You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.05 03:00:01