SparkSQL 之 Parquet 数据转 DataSet 代码实现
·
摘要:Parquet 是 Spark 生态中最推荐的列式存储格式——Schema 自描述(零推断)、谓词下推到 Row Group 级别、列式 Vectorized 批量解码(比逐行快 3~5x)、压缩比高达 3~10x。本文从 spark.read.parquet() 的三种数据源、Parquet 4 级物理布局、谓词下推跳过机制、VectorizedParquetReader、与 JSON/CSV 对比、完整 Options 六个维度,配合 2 张架构图 + 代码实例,全面掌握 Parquet → Dataset 的高效实践。
关键词:Parquet, spark.read.parquet, Predicate Pushdown, Row Group, Vectorized Reader, 列式存储
一、开篇:为什么 Parquet 是标配?
Parquet 四大优势 vs JSON/CSV
① 列式存储: 只读需要的列 → I/O 减少 90%+
② Schema 自描述: Footer 含完整 Schema → 零推断开销
③ 谓词下推: Row Group 级 min/max 过滤 → 跳过不相关数据块
④ 高压缩比: Snappy/Gzip/Zstd → 比文本格式少 3~10x 空间
二、Parquet → Dataset 全流程

2.1 基本读写
// 读取
val df = spark.read.parquet("hdfs://data/events/")
val df = spark.read.parquet("/path/to/file.parquet")
// 写入
df.write.mode("overwrite").parquet("hdfs://output/")
df.write.partitionBy("year", "month").parquet("hdfs://output/")
// Hive Parquet 表
val df = spark.table("hive_db.parquet_table")
2.2 Schema 处理(零推断!)
val df = spark.read.parquet("path") // Schema 自动从 Footer 读取
// 跨文件 Schema 合并
val df = spark.read.option("mergeSchema", "true").parquet("path")
// 分区发现: /data/year=2024/month=01/*.parquet
// → year 和 month 自动成为 DataFrame 列
2.3 DataFrame → Dataset[CaseClass]
case class Event(id: Long, name: String, score: Double)
val ds: Dataset[Event] = spark.read.parquet("hdfs://events/").as[Event]
ds.filter(_.score > 80).map(e => e.copy(score = e.score * 1.1))
三、Parquet 内部结构 & 谓词下推

3.1 4 级物理布局
Parquet File
├── Row Group 0 (~128MB) ← 谓词下推最小粒度
│ ├── Column Chunk: col_a (Dictionary Page + Data Pages)
│ └── Column Chunk: col_b
├── Row Group 1 (~128MB)
└── Footer: Schema + offsets + min/max 统计
3.2 谓词下推
SELECT * FROM events WHERE score > 80
决策过程:
Row Group 0: score.max = 75 → max < 80 → 整组跳过!
Row Group 1: score.max = 95 → 需要读取
Row Group 2: score.min = 85 → 需要读取
结果: 跳过 33% Row Groups,零 I/O!
3.3 Vectorized Reader
传统: for each row → decode col_a → decode col_b → ...
Vectorized: 一次 decode 4096 行 col_a → 4096 行 col_b
优势: 循环展开 · SIMD · 零虚函数 · 缓存友好
性能: 比逐行快 3~5x
四、核心配置
# 读取
spark.sql.parquet.filterPushdown=true
spark.sql.parquet.enableVectorizedReader=true
spark.sql.parquet.mergeSchema=false
# 写入
spark.sql.parquet.compression.codec=snappy # snappy/gzip/lz4/zstd
parquet.block.size=134217728 # 128MB Row Group
parquet.page.size=1048576 # 1MB Page
# 分区
spark.sql.sources.partitionOverwriteMode=static
五、总结
-
Parquet 三大能力:Schema 自描述(零推断)、谓词下推(Row Group 级跳过)、Vectorized 批量解码(3~5x)。
-
4 级物理布局:File → Row Group(128M) → Column Chunk → Page。
-
vs JSON/CSV:列式 + 自描述 + 压缩 3~10x + 零推断。生产环境首选。
作者:starzy
博客:blog.starzy.cn
GitHub:starzy1990.github.io
专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践
更多推荐



所有评论(0)