SparkSQL 之 DataSet 存储代码实现
摘要:DataFrameWriter 提供了统一的写出 API。本文从 6 步流水线、SaveMode 四模式、partitionBy/bucketBy、四大格式对比、saveAsTable/insertInto、文件数控制六个维度,配合 2 张架构图 + 代码实例,全面掌握 Dataset 存储的最佳实践。
关键词:DataFrameWriter, SaveMode, partitionBy, bucketBy, saveAsTable, Parquet, ORC, coalesce
一、开篇
df.write.format("parquet").mode("overwrite").partitionBy("dt").save("path")// 便捷 API: df.write.parquet("path") / df.write.json("path") / df.write.csv("path")二、写出流水线 & SaveMode
SaveMode
| 模式 | 行为 |
|---|---|
| Append | 追加 |
| Overwrite | 覆盖删除旧数据 |
| ErrorIfExists | 目录存在则报错(默认) |
| Ignore | 目录存在则跳过 |
分区与分桶
// 分区: /path/dt=2024-01-01/hour=12/df.write.partitionBy("dt","hour").parquet("path")// 分桶: 10 个桶文件,每个桶内按 ts 排序df.write.bucketBy(10,"user_id").sortBy("ts").saveAsTable("tbl")三、存储格式 & Hive 表写出
格式选型
| 格式 | 特点 |
|---|---|
| Parquet | 默认·列式·Spark原生 |
| ORC | Hive原生·压缩略优 |
| JSON | 可读·体积 3-5x |
| CSV | 通用·体积 5-10x |
saveAsTable vs insertInto
// saveAsTable: 自动建表df.write.mode("overwrite").saveAsTable("db.tbl")// insertInto: 表必须已存在,按列位置匹配!df.write.mode("append").insertInto("db.tbl")四、文件数控制
df.coalesce(4).write.parquet("path")// 无 shuffle → 4 文件df.repartition(4).write.parquet("path")// shuffle → 4 文件五、总结
- 写出流程:write→format→mode→partitionBy→option→save
- 格式:默认 Parquet,Hive 用 ORC,交换用 JSON/CSV
- 技巧:coalesce 减少小文件·bucketBy 加速 JOIN·insertInto 注意列位置
作者:starzy
博客:blog.starzy.cn
GitHub:starzy1990.github.io
专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践