PySpark如何最简最快获取DataFrame分区键并转为Python列表
PySpark快速获取分区键转Python列表方案
方案1:最简原生API写法(仅扫描元数据,无实际数据读取)
Spark 2.3+已原生支持分区列查询优化,针对分区列的去重查询会直接读取目录元数据,不会扫描实际文件内容,执行效率极高,返回结果直接为原生Python列表:
partition_list = ddf2.select("date_str").distinct().rdd.map(lambda x: x[0]).collect()
如果分区值为日期类型,可根据需要直接在map阶段转成字符串格式:
partition_list = ddf2.select("date_str").distinct().rdd.map(lambda x: str(x[0])).collect()
方案2:极致性能方案(无Spark作业提交,Driver端直接提取)
如果需要最高执行效率,无需触发任何Spark作业,可以直接调用Hadoop FileSystem API列取分区目录提取键值,全程仅在Driver端运行:
from py4j.java_gateway import java_import # 获取Hadoop文件系统实例 java_import(spark._jvm, 'org.apache.hadoop.fs.Path') fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(spark._jsc.hadoopConfiguration()) # 遍历分区根目录提取分区值 dir_list = fs.listStatus(spark._jvm.Path("data/bydate")) partition_list = [ d.getPath().getName().split("=")[1] for d in dir_list if d.isDirectory() and "=" in d.getPath().getName() ]
常见问题说明
show partitions指令仅对注册到Hive元数据的分区表生效,直接读取文件生成的临时视图未注册到Hive元数据,因此会报错rdd.getNumPartitions返回的是Spark RDD的内存分区数量,和存储层的分区目录属于完全不同的概念,无法获取存储层分区键
内容的提问来源于stack exchange,提问作者citynorman
相关产品推荐
相关产品推荐

