PySpark如何读取HDFS中指定日期范围的Parquet分区文件?
解决方案:读取HDFS指定日期范围的分区数据
方法1:Spark SQL分区过滤(最推荐)
Spark会自动识别HDFS上的分区列cutoff_date,直接用SQL筛选日期范围即可触发分区裁剪,不会扫描无关分区:
// 加载数据集,Spark自动解析分区列 val df = sqlContext.read.load("/my/path/here") // 注册临时视图 df.createOrReplaceTempView("partitioned_data") // 用字符串范围筛选日期,注意格式和分区一致 val filteredDf = sqlContext.sql("SELECT * FROM partitioned_data WHERE cutoff_date BETWEEN '2023-06-07' AND '2023-10-06'")
方法2:路径通配符直接匹配
如果日期范围可以分段,直接用通配符指定要加载的分区路径,效率最高:
// 组合需要的分区路径 val targetPaths = Seq( // 2023-06-07至2023-06-30 "/my/path/here/cutoff_date=2023-06-0[7-9]", "/my/path/here/cutoff_date=2023-06-[1-3][0-9]", // 2023-07至2023-09全月 "/my/path/here/cutoff_date=2023-0[7-9]-*", // 2023-10-01至2023-10-06 "/my/path/here/cutoff_date=2023-10-0[1-6]" ) // 加载指定路径 val filteredDf = sqlContext.read.load(targetPaths: _*)
方法3:修复load的范围过滤写法
如果之前用sqlContext.read.load的范围写法无效,大概率是没正确利用分区列,试试这两种写法:
// 写法A:先加载再过滤,Spark会自动优化分区裁剪 val filteredDf = sqlContext.read.load("/my/path/here") .filter($"cutoff_date" >= "2023-06-07" && $"cutoff_date" <= "2023-10-06") // 写法B:指定basePath确保分区列被识别 val filteredDf = sqlContext.read .option("basePath", "/my/path/here") .load("/my/path/here/cutoff_date=*") .filter($"cutoff_date".between("2023-06-07", "2023-10-06"))
内容的提问来源于stack exchange,提问作者Arturo Sbr
相关产品推荐
相关产品推荐

