数据读取与保存

Spark 数据读写:文本、JSON、CSV、SequenceFile,什么场景用什么格式

Spark 读写数据,选择什么格式取决于两个因素:

  1. 数据结构:是纯文本、表格、还是键值对?
  2. 性能要求:要快(二进制)还是要通用(文本)?
格式 适用场景 性能 跨语言
文本文件 日志、简单输出 一般 可以
JSON API 数据、半结构化 一般 可以
CSV 表格数据、Excel 导出 一般 可以
SequenceFile Hadoop 生态内部交换 不可以
对象序列化 Spark 内部中间结果 不可以

核心原则:要通用用文本/JSON/CSV,要性能用二进制格式(SequenceFile/Parquet)

文本文件:最通用,但也最基础

什么时候用: 日志文件、简单的文本输出、不需要结构化的数据。

读:

val rdd = sc.textFile("hdfs:///data/logs/2026-07-29.log")

// 指定最小分区数
val rdd = sc.textFile("hdfs:///data/logs/", minPartitions=10)

写:

rdd.saveAsTextFile("hdfs:///output/result")

写入后是一个目录,里面每个分区一个 part-* 文件。

JSON:半结构化的”万能格式”

什么时候用: API 响应、日志中的结构化数据、跨语言数据交换。

有两种方式,推荐用 DataFrame 方式:

方式1:DataFrame 方式(推荐)

val spark = SparkSession.builder().appName("JSON").getOrCreate()

// 读 JSON
val df = spark.read.json("hdfs:///data/people.json")

// 自动推断 schema,可以做 SQL 查询
df.createOrReplaceTempView("people")
spark.sql("SELECT name, age FROM people WHERE age > 18").show()

// 写 JSON
df.write.json("hdfs:///output/json_result")

方式2:RDD 方式(手动解析)

import com.fasterxml.jackson.databind.ObjectMapper

case class Person(name: String, age: Int)
val mapper = new ObjectMapper()

val rdd = sc.textFile("hdfs:///data/people.json")
  .map(line => mapper.readValue(line, classOf[Person]))
  .filter(_ != null)

rdd.map(mapper.writeValueAsString).saveAsTextFile("hdfs:///output/json_rdd")

CSV:表格数据的”标准格式”

什么时候用: Excel 导出、数据库导出的表格数据、数据分析常用的格式。

DataFrame 方式(推荐):

val spark = SparkSession.builder().appName("CSV").getOrCreate()

// 读 CSV
val df = spark.read
  .option("header", "true")        // 第一行是表头
  .option("inferSchema", "true")   // 自动推断数据类型
  .csv("hdfs:///data/users.csv")

// 查
df.show()
df.filter("age > 18").show()

// 写 CSV
df.write
  .option("header", "true")
  .csv("hdfs:///output/csv_result")

常用配置:

配置 说明 示例
header 第一行是不是表头 true / false
inferSchema 自动推断列类型 true / false
delimiter 分隔符(默认逗号) `delimiter=” “`
quote 引号字符 quote="\""
nullValue 空值表示 nullValue="NULL"

SequenceFile:Hadoop 的二进制键值对格式

什么时候用: Hadoop 生态内部数据交换、需要高效存储键值对数据。

import org.apache.hadoop.io.{IntWritable, Text}

// 读 SequenceFile(必须指定 key 和 value 的 Hadoop 类型)
val seqRDD = sc.sequenceFile[Text, IntWritable]("hdfs:///data/seqfile")

// 转成普通类型
val rdd = seqRDD.map { case (k, v) => (k.toString, v.get()) }

// 写 SequenceFile(要转回 Hadoop 类型)
rdd.map { case (k, v) => (new Text(k), new IntWritable(v)) }
  .saveAsSequenceFile("hdfs:///output/seq_result")

注意: SequenceFile 是 Hadoop 专属格式,跨语言兼容性差。如果你的数据不需要被 Hadoop 外的系统读,用 SequenceFile 性能很好。

对象序列化:Spark 内部的”快速缓存”

什么时候用: 在 Spark 作业之间传递中间结果,不需要跨语言读。

// 定义可序列化的类
case class User(id: Long, name: String) extends Serializable

val users = sc.parallelize(Seq(User(1, "Alice"), User(2, "Bob")))

// 保存为对象文件
users.saveAsObjectFile("hdfs:///output/objects")

// 读取
val loaded = sc.objectFile[User]("hdfs:///output/objects")

性能优化: 默认用 Java 序列化,速度慢。换成 Kryo 序列化:

val conf = new SparkConf()
  .set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
  .registerKryoClasses(Array(classOf[User]))

高级优化技巧

1. 输出压缩

// 压缩 CSV 输出
df.write
  .option("header", "true")
  .option("compression", "snappy")
  .csv("hdfs:///output/compressed")

支持的压缩算法:snappygziplzobzip2

2. 分区写入

// 按日期分区写入
df.write
  .partitionBy("date")
  .parquet("hdfs:///output/partitioned")

3. 控制输出文件数量

// 减少分区数,避免生成太多小文件
df.coalesce(10).write.csv("hdfs:///output/fewer_files")

// 增加分区数,提高并行度
df.repartition(50).write.csv("hdfs:///output/more_parallel")