SparkSQL 之 JDBC 数据转 DataSet 代码实现
·
摘要:JDBC 是连接传统关系型数据库的桥梁。本文从 JDBC 读取三模式(整表/数值分区/自定义Predicate)、并行分区原理、谓词/列裁剪下推、批量写入、连接池管理、以及四大常见坑六个维度,配合 2 张架构图 + 完整代码实例,覆盖 JDBC 操作的全部实践要点。
关键词:spark.read.jdbc, JDBC, partitionColumn, numPartitions, Predicate Pushdown, batchsize
一、开篇
Spark 通过 JDBC 连接所有标准 JDBC 兼容数据库,核心 API 就是 spark.read.jdbc()。
val props = new java.util.Properties()
props.setProperty("user", "root")
props.setProperty("password", "123456")
props.setProperty("driver", "com.mysql.cj.jdbc.Driver")
val url = "jdbc:mysql://host:3306/db"
val df = spark.read.jdbc(url, "users", props)
二、JDBC 读取全流程

2.1 三种读取入口
// 方式 1: spark.read.jdbc
val df = spark.read.jdbc(url, "users", props)
// 方式 2: format("jdbc").options()
val df = spark.read.format("jdbc").option(...).load()
// 方式 3: 子查询
val df = spark.read.jdbc(url,
"(SELECT id,name FROM users WHERE status=1) AS u", props)
2.2 并行分区读取
// 数值列等分区间
val df = spark.read.format("jdbc")
.option("partitionColumn", "id")
.option("lowerBound", "1")
.option("upperBound", "10000000")
.option("numPartitions", "20")
.load()
// → 20 个 Task,每个执行一个 WHERE id BETWEEN ... AND ...
// 自定义 Predicate 列表
val predicates = Array("gender = 'M'", "gender = 'F'")
val df = spark.read.jdbc(url, "users", predicates, props)
2.3 DataFrame → Dataset[CaseClass]
case class User(id: Long, name: String, age: Int)
val ds: Dataset[User] = spark.read.jdbc(url, "users", props).as[User]
三、连接管理 & 完整代码模式

3.1 谓词/列裁剪下推
spark.read.jdbc(url, "users", props)
.filter("age > 30")
.select("id", "name", "age")
// → SQL: SELECT id, name, age FROM users WHERE age > 30
3.2 批量写回
df.write.mode("append")
.option("batchsize", "5000")
.option("isolationLevel", "READ_UNCOMMITTED")
.jdbc(url, "target_table", props)
四、四大常见坑
① 连接数爆炸: numPartitions × executors 个连接 → DB max_connections 必须足够
② 数据倾斜: 分区列值分布不均 → 长尾 Task → 用自定义 Predicate 解决
③ 全量拉取: 未加 filter → 全表扫描 → 读时用 query 限定范围
④ batchsize 太小: 默认 1000 → 增量到 5000~10000 显著提速
五、总结
- 三种读取模式:整表/subquery + 数值列分区 + 自定义 Predicate。推荐用分区并行读。
- 优化要点:谓词/列裁剪下推 + numPartitions ≤ 20 + batchsize=5000~10000。
- 避坑:控制连接数、避免分区倾斜、查询加 WHERE 限制。
作者:starzy
博客:blog.starzy.cn
GitHub:starzy1990.github.io
专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践
更多推荐



所有评论(0)