spark连接jdbc
Spark 连接 JDBC:读 MySQL 用 DataFrame,别用 JdbcRDD 了
Spark 连接关系型数据库(MySQL、PostgreSQL 等),有两种方式:
- JdbcRDD 只能读,不能写,代码繁琐。早期 RDD API,已过时
- DataFrame API 读写都支持,语法简洁,性能更好
准备工作:加驱动、配连接
1. 加 JDBC 驱动
<!-- Maven -->
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
<version>8.0.28</version>
</dependency>
如果用 spark-submit 提交:
spark-submit --jars mysql-connector-java-8.0.28.jar your-app.jar
2. 连接参数
val jdbcUrl = "jdbc:mysql://localhost:3306/test?useSSL=false"
val props = new java.util.Properties()
props.setProperty("user", "root")
props.setProperty("password", "123456")
props.setProperty("driver", "com.mysql.cj.jdbc.Driver")
MySQL 驱动类名:
| MySQL 版本 | 驱动类 |
|---|---|
| 5.x | com.mysql.jdbc.Driver |
| 8.x | com.mysql.cj.jdbc.Driver |
DataFrame
读数据:从 MySQL 到 DataFrame
方式1:读整张表
val df = spark.read.jdbc(jdbcUrl, "user", props)
df.show()
方式2:并行读取(大表推荐)
val df = spark.read.jdbc(
jdbcUrl,
"user",
columnName = "id", // 分区列(数值型)
lowerBound = 1, // 下界
upperBound = 10000, // 上界
numPartitions = 10, // 分区数
props
)
这样 10 个 Task 并行读,每个 Task 读 id 的一段范围(比如 1-1000,1001-2000…)。
方式3:用 SQL 查(只读需要的行/列)
val df = spark.read.jdbc(
jdbcUrl,
"(SELECT id, name FROM user WHERE age > 18) AS sub",
props
)
写数据:从 DataFrame 到 MySQL
// 准备数据
val data = Seq(
(101, "Alice", 25),
(102, "Bob", 30)
).toDF("id", "name", "age")
// 写入 MySQL
data.write.jdbc(jdbcUrl, "user", props)
写入模式(.mode()):
| 模式 | 行为 |
|---|---|
append |
追加数据(默认) |
overwrite |
覆盖表(先删后建) |
ignore |
表存在就忽略,不写 |
error |
表存在就抛异常 |
data.write
.mode("overwrite")
.jdbc(jdbcUrl, "user", props)
批量写入优化:
props.setProperty("batchsize", "1000") // 每批 1000 条
props.setProperty("isolationLevel", "NONE") // 关事务,提升速度
data.write.jdbc(jdbcUrl, "user", props)
JdbcRDD(了解即可,不推荐)
JdbcRDD 是旧 API,只能读,不能写,代码还长:
import org.apache.spark.rdd.JdbcRDD
val rdd = new JdbcRDD(
sc,
() => DriverManager.getConnection(jdbcUrl, props),
"SELECT id, name FROM user WHERE id >= ? AND id <= ?",
1, 100, 10, // 下界、上界、分区数
rs => (rs.getInt("id"), rs.getString("name"))
)
JdbcRDD 能做的,DataFrame 都能做,还更好。别用 JdbcRDD 了。
常见问题
1. ClassNotFoundException: com.mysql.cj.jdbc.Driver
驱动没加。确认:
pom.xml有依赖spark-submit带了--jars- IDE 本地运行的话,驱动 JAR 在 classpath 里
2. 连接超时
报 CommunicationsException: Communications link failure。
加超时参数:
val jdbcUrl = "jdbc:mysql://localhost:3306/test?connectTimeout=10000&socketTimeout=30000"
3. 并行读取时某些分区慢(数据倾斜)
分区列数据分布不均(比如 id 不是均匀递增的)。
换分区列,或者手动算 lowerBound 和 upperBound 的范围。
4. 写入报权限不足
MySQL 用户没有 INSERT 或 CREATE 权限。
GRANT INSERT, CREATE ON test.* TO 'root'@'localhost';
FLUSH PRIVILEGES;