news 2026/8/19 13:00:19

SparkSQL 之 JDBC 数据转 DataSet 代码实现

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
SparkSQL 之 JDBC 数据转 DataSet 代码实现

摘要:JDBC 是连接传统关系型数据库的桥梁。本文从 JDBC 读取三模式(整表/数值分区/自定义Predicate)、并行分区原理、谓词/列裁剪下推、批量写入、连接池管理、以及四大常见坑六个维度,配合 2 张架构图 + 完整代码实例,覆盖 JDBC 操作的全部实践要点。

关键词:spark.read.jdbc, JDBC, partitionColumn, numPartitions, Predicate Pushdown, batchsize


一、开篇

Spark 通过 JDBC 连接所有标准 JDBC 兼容数据库,核心 API 就是spark.read.jdbc()

valprops=newjava.util.Properties()props.setProperty("user","root")props.setProperty("password","123456")props.setProperty("driver","com.mysql.cj.jdbc.Driver")valurl="jdbc:mysql://host:3306/db"valdf=spark.read.jdbc(url,"users",props)

二、JDBC 读取全流程

2.1 三种读取入口

// 方式 1: spark.read.jdbcvaldf=spark.read.jdbc(url,"users",props)// 方式 2: format("jdbc").options()valdf=spark.read.format("jdbc").option(...).load()// 方式 3: 子查询valdf=spark.read.jdbc(url,"(SELECT id,name FROM users WHERE status=1) AS u",props)

2.2 并行分区读取

// 数值列等分区间valdf=spark.read.format("jdbc").option("partitionColumn","id").option("lowerBound","1").option("upperBound","10000000").option("numPartitions","20").load()// → 20 个 Task,每个执行一个 WHERE id BETWEEN ... AND ...// 自定义 Predicate 列表valpredicates=Array("gender = 'M'","gender = 'F'")valdf=spark.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 > 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 显著提速

五、总结

  1. 三种读取模式:整表/subquery + 数值列分区 + 自定义 Predicate。推荐用分区并行读。
  2. 优化要点:谓词/列裁剪下推 + numPartitions ≤ 20 + batchsize=5000~10000。
  3. 避坑:控制连接数、避免分区倾斜、查询加 WHERE 限制。

作者:starzy
博客:blog.starzy.cn
GitHub:starzy1990.github.io
专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/19 13:00:00

基于大语言模型的群体推荐系统:AgentGR模拟器原理与实践

1. 从“众口难调”到“智能调和”:群体推荐的困境与AgentGR的破局在推荐系统的世界里,给单个用户推荐内容已经是一门相当成熟的学问。从协同过滤到深度学习,算法工程师们构建了无数模型,试图精准预测“你”的下一次点击。然而&…

作者头像 李华
网站建设 2026/8/19 12:59:18

奔驰概念雕塑:解码软件定义汽车时代的数字化转型与设计变革

1. 从“三叉星徽”到“数字星云”:奔驰概念雕塑的叙事转向 最近,奔驰在各大设计周和前瞻技术展上亮相的一系列概念雕塑,成了圈内人热议的话题。如果你还认为这些只是造型更酷炫、线条更犀利的“未来汽车模型”,那可能就错过了它背…

作者头像 李华
网站建设 2026/8/19 12:58:34

Python爬虫实战:招聘网站职位信息采集系统(完整版)

一、项目背景与需求分析 在当今数字化时代,招聘数据已成为洞察就业市场、分析行业趋势、指导职业规划的重要资源。无论是求职者寻找合适岗位、HR调研薪酬水平,还是数据分析师研究就业市场动态,能够高效、准确地采集招聘网站职位信息都是一项极具价值的技能。本篇文章将带领…

作者头像 李华
网站建设 2026/8/19 12:57:52

第 6 篇:「用数据说话」— 性能基准测试如何证明架构

开场:感觉不等于事实 到目前为止,我们讨论了很多次"快速"“性能好”。 但有一个问题:你怎么知道真的快? 感觉骗人。优化前"好像"卡,优化后"好像"快。直到你看到数据。 这一篇讲的是 LoopAgent 如何用科学的基准测试来证明架构的有效性。…

作者头像 李华
网站建设 2026/8/19 12:57:46

开源无人机DIY全攻略:从STM32飞控到Betaflight调参实战

1. 项目概述:当开源精神遇上DIY乐趣 如果你和我一样,既着迷于无人机在天空自由翱翔的姿态,又对“黑匣子”般的商业产品内部构造充满好奇,总想亲手拆解、改造,那么“Stamp Fly”这个名字,你一定会感兴趣。这…

作者头像 李华
网站建设 2026/8/19 12:55:27

DeepSeek Harness 上手指南

文章目录前言一、认识 DeepSeek Harness:它是什么1.1 一切皆插件1.2 先记几个名词1.3 它能做什么二、安装:两分钟跑起来2.1 环境要求2.2 Windows 上执行 npm -v 报错怎么办2.3 npm 和 npx 的区别2.4 方式一:临时运行2.5 方式二:全…

作者头像 李华