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 不是均匀递增的)。

换分区列,或者手动算 lowerBoundupperBound 的范围。

4. 写入报权限不足

MySQL 用户没有 INSERT 或 CREATE 权限。

GRANT INSERT, CREATE ON test.* TO 'root'@'localhost';
FLUSH PRIVILEGES;