如何使用Scala语言读取HDFS中的CSV数据集?
嘿,这个问题我太熟了!用Scala读取HDFS上的CSV文件,其实有几种实用的方式,我给你拆解清楚,从底层操作到大数据工具都覆盖到:
方法一:使用Hadoop原生API(适合底层精细控制)
如果需要对文件读取过程做精细控制(比如逐行处理、自定义编码),直接用Hadoop的FileSystem API是个不错的选择。这里给你一个完整的示例,包括资源的安全关闭:
import org.apache.hadoop.fs.{FileSystem, Path} import org.apache.hadoop.conf.Configuration import java.io.BufferedReader import java.io.InputStreamReader object HdfsCsvReader { def main(args: Array[String]): Unit = { // 初始化Hadoop配置,自动加载classpath下的core-site.xml、hdfs-site.xml val conf = new Configuration() // 如果你的Hadoop集群没有配置默认FS,手动指定:conf.set("fs.defaultFS", "hdfs://your-namenode:9000") // 获取FileSystem实例 val fs = FileSystem.get(conf) val csvPath = new Path("/path/to/your/file.csv") // 打开文件流并读取 val inputStream = fs.open(csvPath) val reader = new BufferedReader(new InputStreamReader(inputStream)) try { var line: String = null // 跳过表头(如果需要的话) line = reader.readLine() while ({line = reader.readLine(); line != null}) { val fields = line.split(",") // 根据CSV的分隔符调整,比如分号就用";" // 这里处理每一行数据,比如打印字段 println(s"字段1: ${fields(0)}, 字段2: ${fields(1)}") } } finally { // 务必关闭资源,避免内存泄漏 reader.close() inputStream.close() fs.close() } } }
小提示:如果用的是Scala 2.13及以上版本,可以用Using语法来自动管理资源,不用手动写finally:
import scala.util.Using Using.Manager { use => val inputStream = use(fs.open(csvPath)) val reader = use(new BufferedReader(new InputStreamReader(inputStream))) // 读取逻辑同上 } match { case Right(_) => println("读取完成") case Left(e) => e.printStackTrace() }
方法二:使用Apache Spark(大数据场景首选)
如果你的场景是处理数据集(哪怕是小数据集),Spark绝对是最省心的选择——它自带CSV解析器,能自动处理表头、类型推断,还能快速转换成DataFrame/Dataset进行后续操作。
方式A:DataFrame API(推荐)
import org.apache.spark.sql.SparkSession object SparkCsvReader { def main(args: Array[String]): Unit = { // 初始化SparkSession val spark = SparkSession.builder() .appName("ReadHdfsCsv") .master("local[*]") // 本地调试用,生产环境去掉这行 .getOrCreate() // 读取HDFS上的CSV文件 val df = spark.read .option("header", "true") // 如果CSV有表头,设为true .option("inferSchema", "true") // 自动推断字段类型 .csv("hdfs://your-namenode:9000/path/to/your/file.csv") // 查看前几行数据 df.show() // 如果要转换成Scala集合(因为数据集不大) val dataList = df.collect().toList dataList.foreach(row => println(row.mkString(", "))) // 关闭SparkSession spark.stop() } }
方式B:RDD API(更底层的弹性分布式数据集)
如果需要更灵活的行级处理,可以用RDD:
import org.apache.spark.SparkContext import org.apache.spark.SparkConf object RddCsvReader { def main(args: Array[String]): Unit = { val conf = new SparkConf().setAppName("ReadHdfsCsvRdd").setMaster("local[*]") val sc = new SparkContext(conf) // 读取文件为RDD[String] val csvRdd = sc.textFile("hdfs://your-namenode:9000/path/to/your/file.csv") // 跳过表头,分割每行字段 val dataRdd = csvRdd.zipWithIndex().filter(_._2 > 0).map(_._1.split(",")) // 处理数据 dataRdd.foreach(fields => println(s"字段: ${fields.mkString("|")}")) sc.stop() } }
关键注意事项
- HDFS路径格式:必须是
hdfs://namenode-host:port/文件路径,如果集群配置了默认FS,也可以直接用/文件路径 - 依赖配置:如果用sbt构建项目,记得添加对应的依赖:
- Hadoop API:
libraryDependencies += "org.apache.hadoop" % "hadoop-common" % "3.3.4"(版本根据你的Hadoop集群调整) - Spark:
libraryDependencies += "org.apache.spark" %% "spark-sql" % "3.4.0"(Scala版本要和Spark对应)
- Hadoop API:
- 特殊CSV格式:如果CSV包含引号、转义字符或者自定义分隔符,记得在读取时添加对应的option,比如Spark里的
.option("quote", "\"")、.option("delimiter", ";")
内容的提问来源于stack exchange,提问作者Sunitha
相关产品推荐
相关产品推荐

