Spark Python如何在读取时创建DataFrame分区及查询分区数
Spark DataFrame 分区相关问题解答
1. 读取数据生成DataFrame时可直接创建分区
不同数据源的分区配置规则如下:
- 文件类数据源(CSV、JSON、Parquet、普通文本等):Spark读取时默认按照文件块大小、文件数量自动生成分区,也可以通过参数手动控制分区:
- 调整
spark.sql.files.maxPartitionBytes参数设置单个分区的最大存储阈值,间接控制总分区数 - 大量小文件场景可开启小文件合并配置,部分数据源支持直接通过
option("numPartitions", 目标分区数)在读取阶段指定分区数
- 调整
- JDBC类数据源:读取时指定分区参数即可直接生成多分区,示例代码如下:
# Python 示例:JDBC读取时直接指定10个分区 df = spark.read.format("jdbc") \ .option("url", "jdbc:mysql://你的数据库地址/库名") \ .option("dbtable", "表名") \ .option("user", "数据库账号") \ .option("password", "数据库密码") \ .option("partitionColumn", "用作分区的整数列名,比如id") \ .option("lowerBound", 分区列的最小值) \ .option("upperBound", 分区列的最大值) \ .option("numPartitions", 10) \ .load()
2. DataFrame分区数量查询方法
你使用的df.rdd.getNumPartitions()是标准正确的查询方法,默认返回1通常由以下场景导致:
- 读取的是单个小于默认分区阈值(通常为128M)的不可拆分小文件,Spark不会自动拆分生成分区
- Spark运行在本地单核心模式,默认生成1个分区
- 读取后执行过合并分区的算子(比如
coalesce(1))、全局排序等会合并分区的操作
coalesce()与repartition()调整分区的注意事项
repartition(目标分区数):会触发全量数据shuffle,支持调大或调小分区数,调整后分区数据分布均匀,也支持按指定字段哈希分区(用法为repartition(目标分区数, 列名)),适合分区数大幅调整、需要按字段分区的场景coalesce(目标分区数):仅支持调小分区数,不会触发全量shuffle,通过合并相邻分区实现,性能远高于repartition(),适合减少分区数的场景;如果强行用coalesce()调大分区数,操作不会生效,分区数会保持原有数量
内容的提问来源于stack exchange,提问作者Mubeen
相关产品推荐
相关产品推荐

