ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

SparkSQL 之 JDBC 数据转 DataSet 代码实现

SparkSQL 之 JDBC 数据转 DataSet 代码实现 摘要JDBC 是连接传统关系型数据库的桥梁。本文从 JDBC 读取三模式整表/数值分区/自定义Predicate、并行分区原理、谓词/列裁剪下推、批量写入、连接池管理、以及四大常见坑六个维度配合 2 张架构图 完整代码实例覆盖 JDBC 操作的全部实践要点。关键词spark.read.jdbc, JDBC, partitionColumn, numPartitions, Predicate Pushdown, batchsize一、开篇Spark 通过 JDBC 连接所有标准 JDBC 兼容数据库核心 API 就是spark.read.jdbc()。valpropsnewjava.util.Properties()props.setProperty(user,root)props.setProperty(password,123456)props.setProperty(driver,com.mysql.cj.jdbc.Driver)valurljdbc:mysql://host:3306/dbvaldfspark.read.jdbc(url,users,props)二、JDBC 读取全流程2.1 三种读取入口// 方式 1: spark.read.jdbcvaldfspark.read.jdbc(url,users,props)// 方式 2: format(jdbc).options()valdfspark.read.format(jdbc).option(...).load()// 方式 3: 子查询valdfspark.read.jdbc(url,(SELECT id,name FROM users WHERE status1) AS u,props)2.2 并行分区读取// 数值列等分区间valdfspark.read.format(jdbc).option(partitionColumn,id).option(lowerBound,1).option(upperBound,10000000).option(numPartitions,20).load()// → 20 个 Task每个执行一个 WHERE id BETWEEN ... AND ...// 自定义 Predicate 列表valpredicatesArray(gender M,gender F)valdfspark.read.jdbc(url,users,predicates,props)2.3 DataFrame → Dataset[CaseClass]caseclassUser(id:Long,name:String,age:Int)valds: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 303.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 batchsize5000~10000。避坑控制连接数、避免分区倾斜、查询加 WHERE 限制。作者starzy博客blog.starzy.cnGitHubstarzy1990.github.io专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践
返回列表