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

Scala如何从HDFS目录获取最大load分区并读取其中文件

自动获取HDFS中最大load值分区并读取CSV文件(Scala实现)

方法一:直接操作HDFS目录(推荐,效率更高)

通过HDFS FileSystem API遍历目录、提取分区值,精准定位最大load分区后再读取文件:

import org.apache.hadoop.fs.{FileSystem, Path}
import org.apache.spark.sql.SparkSession

// 初始化SparkSession
val spark = SparkSession.builder()
  .appName("ReadMaxLoadPartition")
  .getOrCreate()

// 获取HDFS文件系统实例
val fs = FileSystem.get(spark.sparkContext.hadoopConfiguration)
val baseDir = new Path("hdfs://device/signs/")

// 筛选出所有load前缀的分区目录
val loadPartitions = fs.listStatus(baseDir)
  .filter(_.isDirectory)
  .map(_.getPath.getName)
  .filter(_.startsWith("load="))

// 提取load数值并找到最大值(用maxOption避免空分区报错)
val maxLoad = loadPartitions
  .map(partition => partition.split("=")(1).toInt)
  .maxOption

// 读取最大load分区的CSV文件
val read_csv = maxLoad match {
  case Some(loadVal) => 
    spark.read.format("csv")
      .load(s"hdfs://device/signs/load=$loadVal")
  case None => 
    // 无有效分区时返回空DataFrame,可根据需求调整逻辑
    spark.emptyDataFrame
}

关键逻辑说明:

  • FileSystem.get:复用Spark的Hadoop配置获取HDFS连接
  • 先过滤目录、再提取load=后的数值,避免无效文件干扰
  • maxOption处理无有效分区的边界情况,防止代码崩溃

方法二:利用Spark分区自动发现(简单但效率低)

如果分区数量较少,可直接读取所有load分区,再通过字段过滤最大值:

import org.apache.spark.sql.functions.max

// 读取所有load分区,指定basePath让Spark识别分区字段
val allPartitionsDF = spark.read.format("csv")
  .option("basePath", "hdfs://device/signs/")
  .load("hdfs://device/signs/load=*")

// 获取最大load值
val maxLoadVal = allPartitionsDF.select("load").distinct()
  .agg(max("load"))
  .head()
  .getInt(0)

// 过滤出最大load分区的数据
val read_csv = allPartitionsDF.filter(s"load = $maxLoadVal")

注意:

  • 此方法会先加载所有分区的元数据,分区数量多的时候性能不如方法一
  • 需确保Spark能自动识别load为分区字段(无需在CSV文件中包含该列)

内容的提问来源于stack exchange,提问作者Pluto

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 05:05:27