PySpark读取同一份CSV时分区数不稳定的原因咨询
Spark读取CSV分区数波动的原因解析
一、Spark读取CSV的默认分区数怎么算?
Spark读取CSV这类文本文件时,默认分区数主要由以下几个因素决定:
spark.sql.files.maxPartitionBytes:默认值128MB,是单个分区允许存储的最大字节数spark.sql.files.openCostInBytes:默认值4MB,用来估算打开文件的IO成本,计算分区数时会把这个值加到文件总大小里- 文件的物理分片:如果文件存储在HDFS这类分块文件系统上,HDFS的块大小(默认128MB)也会影响初始分区的划分
简化计算公式为:
分区数 = 向上取整( (文件总大小 + 打开文件的IO成本) / 单个分区最大字节数 )
你的175.46MB CSV文件按默认参数计算,理论分区数应为2,但实际出现16的情况,说明你的文件在存储层被拆分成了多个小分片(比如HDFS写入时的分片、或是由多个小文件合并成的逻辑大文件)。
二、为什么分区数会波动且逐渐减少?
你碰到的重复读取时分区数无规律减少、且绝不超过之前数值的现象,主要和YARN集群的资源调度以及Spark的缓存机制有关:
- YARN容器动态变化
YARN模式下,Executor容器是动态分配的。第一次读取文件时,集群资源充足,Spark会根据文件分片数和可用Executor数量,拆分出最多的分区(比如16个)。后续读取时,部分Executor可能因集群资源紧张被回收、或节点负载过高,Spark会自动合并分区以适配当前可用计算资源,导致分区数减少。 - Spark文件缓存优化
第一次读取后,Spark会将部分或全部数据块缓存到Executor的内存或磁盘中。后续读取时会优先读取缓存数据,而缓存的块可能被合并存储(比如多个小分区合并为一个大缓存块),从而减少了后续读取的分区数。 - IO性能自适应调整
Spark的文件读取器会根据实际IO性能动态调整分区策略。第一次读取若遇到部分节点IO延迟高,Spark会拆分更多分区并行处理;后续读取IO性能稳定后,会合并不必要的小分区以降低调度开销。
三、如何固定分区数?
如果需要每次读取的分区数一致,可以在读取时显式指定:
# 方法1:读取后重分区 df = spark.read.csv("你的文件路径", header=True, inferSchema=True).repartition(8) # 方法2:读取时直接指定分区数 df = spark.read.option("numPartitions", 8).csv("你的文件路径")
内容的提问来源于stack exchange,提问作者figs_and_nuts
相关产品推荐
相关产品推荐

