数据读取与保存
Spark 数据读写:文本、JSON、CSV、SequenceFile,什么场景用什么格式
Spark 读写数据,选择什么格式取决于两个因素:
- 数据结构:是纯文本、表格、还是键值对?
- 性能要求:要快(二进制)还是要通用(文本)?
| 格式 | 适用场景 | 性能 | 跨语言 |
|---|---|---|---|
| 文本文件 | 日志、简单输出 | 一般 | 可以 |
| 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")
支持的压缩算法:snappy、gzip、lzo、bzip2。
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")