ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Spark连接达梦数据库全攻略:JDBC驱动、方言注册与性能调优实战

Spark连接达梦数据库全攻略:JDBC驱动、方言注册与性能调优实战 最近在搞数仓和数据集成的项目遇到一个很典型的场景计算引擎是Spark底下的业务库和归档库却是达梦数据库DM8。一开始我以为就是把连接串从MySQL换成达梦而已结果真正联调起来才发现Spark对着达梦还真不是开箱即用驱动要处理、方言要注册、甚至一些常见的数据类型都得分外小心。折腾了几天踩了不少坑把整个过程整理出来希望对正在做类似国产化替换或者数据平台对接的朋友有帮助。这篇文章会把Spark适配达梦从原理到实操讲清楚。内容包括为什么Spark不能直接连达梦、JDBC驱动怎么部署、读取和写入的完整示例、如何注册自定义方言解决SQL语法不兼容以及我在生产环境里遇到过的问题和排查方法。适合正在做大数据平台对接达梦的工程师也适合刚接触国产数据库适配的同学参考。1. 适配原理与方案选型1.1 Spark访问达梦的核心链路先说结论Spark本身没有达梦原生的连接器目前最主流、最省事的做法是通过JDBC桥接。可能有人会问达梦不是号称高度兼容Oracle吗那直接用Spark的Oracle方言行不行我实测下来这条路不算完全走不通但非常不推荐。原因在于Spark识别数据库方言不是靠你自己声称“我连的是Oracle”而是通过JdbcDialects.get(url)去匹配。达梦的JDBC URL格式是jdbc:dm://host:5236Spark内置的方言表里没有这种前缀所以它只会给你一个默认的NoopDialect导致后续生成的SQL在引号、分页、类型映射这些细节上统统不对。整个访问链路是这样一个过程Spark通过用户提供的url、user、password、driver参数创建连接然后会向达梦下推一些查询比如读取表结构、获取列信息、计算统计信息最后真正拉数据时会生成类似SELECT * FROM (SELECT ... FROM 表名 WHERE 条件) AS tmp的嵌套子查询交给达梦执行。这个过程中任何一环出的不兼容都会直接报错。所以在动手之前我对这个适配方案做了个判断与其去“骗”Spark认为达梦是Oracle不如老老实实告诉它这是一个新的数据库种类然后帮它把这个数据库的方言注册进去。这也是后面章节里最核心的操作。1.2 三种常见适配方案的取舍我在调研阶段其实列了三个备选方案第一是直接JDBC读写。就是Spark把达梦当成一个普通的JDBC数据源通过format(jdbc)去读去写。优点是实现成本低、代码改动少缺点是对达梦的特定类型和SQL方言兼容不够需要额外注册方言来弥补而且大批量写入时性能一般。第二是中间库中转。让Spark先把数据写到MySQL或者PostgreSQL再通过达梦的DBLINK或者第三方同步工具导入达梦。优点是能绕开Spark对达梦方言支持不足的问题缺点是链路变长实时性差而且等于没做真正的适配运维的时候多了一个中间环节很别扭。第三是自定义DataSource。基于Spark的数据源V2接口为达梦写一个完整的连接器。这个方案最彻底功能、性能、类型映射都能做到最优但开发量大而且要对Spark的DataSource API非常熟。对于大多数业务场景来说属于杀鸡用牛刀。我最终选了方案一并为达梦注册一个自定义方言。理由很简单这个方案能覆盖90%以上的日常需求改动量小团队后续接手也容易踩坑之后把经验沉淀成公共模块一劳永逸。2. 环境准备与依赖配置2.1 达梦JDBC驱动的准备达梦的JDBC驱动一般是一份DmJdbcDriver18.jar在达梦数据库安装目录的drivers/jdbc下面能直接找到。值得注意的是驱动包的版本和达梦库的版本最好保持一致跨大版本容易出现一些奇怪的连接报错。拿到jar包之后需要把它分发到所有运行Spark任务的节点上。如果是本地测试用spark-shell --jars /path/DmJdbcDriver18.jar或者spark-submit --jars就行。如果是跑在YARN集群上单靠--jars有时候会因为节点本地没有驱动而出现各种ClassNotFound更稳妥的做法是把这个jar包放到Spark的jars目录下或者用--driver-class-path和--executor-extra-class-path分别指定清楚。驱动类名这里必须强调一下很多人一看包名是DmJdbcDriver18就以为driver参数要写成DmJdbcDriver18。实际上正确的全限定类名是dm.jdbc.driver.DmDriver这个类名拼错是新手最容易踩的坑报错信息通常会显示ClassNotFoundException: dm.jdbc.driver.DmDriver。另外连接URL的格式是jdbc:dm://IP:5236默认端口是5236如果达梦实例改过端口要跟着改。还有一点如果达梦实例开了数据库模式比如MySQL兼容模式URL末尾可以追加参数比如jdbc:dm://IP:5236?compatibleModemysql但我在生产环境里一般不太依赖这个毕竟SQL方言的问题交给后面自定义方言去处理更可控。2.2 连接参数与Spark环境配置除了驱动和URL还有一组参数需要提前配置好不然Spark可能连上数据库却拿不到正确的连接属性。我在项目里一般会配这么几项user和password达梦用户的账号密码。fetchsize每次从数据库拉取的记录数建议设成500或1000。不设的话达梦驱动默认可能会一次性把所有结果集拉进来对Spark的Driver内存压力很大。batchsize写入批次大小建议1000~5000需要根据字段宽度和网络带宽去压并不是越大越好。isolationLevel事务隔离级别。写入若遇到锁冲突可以试试READ_UNCOMMITTED但注意不要无脑调低隔离级别要结合业务场景判断。socketTimeout和connectTimeout默认不配可能会因为网络抖动导致连接长期挂起。我一般给socketTimeout设600秒connectTimeout设30秒防止某个Executor拉数据时卡死把整个任务拖垮。在Spark侧如果要让所有作业都能连达梦最好在spark-defaults.conf里加上一行spark.jars.packages或spark.jars指向驱动路径避免每次提交都带参数。代码里如果用的是spark.read.format(jdbc)驱动类名和URL直接写死在option里即可。val jdbcDF spark.read .format(jdbc) .option(driver, dm.jdbc.driver.DmDriver) .option(url, jdbc:dm://127.0.0.1:5236) .option(user, SYSDBA) .option(password, ******) .option(dbtable, TEST_USER) .option(fetchsize, 1000) .load()这里注意一点dbtable不一定非得是表名也可以是一个子查询比如(select * from test_user where age 20) tmpSpark会把子查询包装成一个临时的数据集。但是子查询在这层包装之后容易和达梦的语法起冲突所以能用表名尽量用表名复杂过滤条件放到Spark侧或者放到后面的自定义查询里。3. Spark读写达梦的完整实操3.1 DataFrame读取达梦数据参数与SQL下推基本读取方式上面已经给了示例但生产环境里不会只满足于把整张表拉出来。如果达梦表有几千万甚至上亿行直接全量读会产生两个问题一是Driver端拉取太多元数据内存爆掉二是单分区读取整个Spark任务的并行度完全没发挥出来。解决方案是让Spark对达梦表做分区读取。JDBC数据源支持用partitionColumn、lowerBound、upperBound、numPartitions这四个参数组合来实现多个Executor并行去查数据库。原理是Spark会把这四个参数拼成一个范围条件比如SELECT * FROM test_user WHERE id 0 AND id 10000 SELECT * FROM test_user WHERE id 10000 AND id 20000这类条件分别下推给不同的Executor去执行相当于多个连接同时去达梦拉数据。选择分区键的时候有几个硬性要求列必须是数值型且最好是主键或者有索引否则每个子查询都会做全表扫描性能反而更差。val jdbcDF spark.read .format(jdbc) .option(driver, dm.jdbc.driver.DmDriver) .option(url, jdbc:dm://127.0.0.1:5236) .option(user, SYSDBA) .option(password, ******) .option(dbtable, TEST_USER) .option(partitionColumn, id) .option(lowerBound, 1) .option(upperBound, 10000000) .option(numPartitions, 10) .option(fetchsize, 1000) .load()这里有个细节容易踩坑partitionColumn指定的列必须能在SQL里直接用于大小比较且取值范围最好覆盖整张表的真实数据范围。如果lowerBound设成0但表里最小id是10000也没关系最多是第一个子查询扫出来为空不会报错。但如果numPartitions设得过大而这张表的索引选择性又不好达梦那边会同时涌进好多大查询直接把数据库连接池打满任务倒是没报错就是慢得离谱。3.2 复杂查询与表连接Spark SQL接达梦在实际项目中我很少直接对着达梦表全量处理更多场景是达梦里已经有了一张业务宽表或者指标表Spark这边要把这表和其他数据源比如Hive、另一个MySQL库做关联加工。这时候最简单的办法是把达梦表注册成临时视图jdbcDF.createOrReplaceTempView(dm_test_user) val result spark.sql( SELECT a.id, a.name, b.order_id FROM dm_test_user a JOIN hive_orders b ON a.id b.user_id )Spark SQL支持对临时视图做join、过滤、聚合这些操作会尽量下推到JDBC数据源也就是达梦先去算一部分再交回Spark做后续逻辑。但要注意下推的力度取决于你对SQL的写法。比如SELECT * FROM dm_test_user WHERE age 20这种Spark会把age 20下推到达梦但如果你在Spark侧用了自定义UDF那就铁定是先把全表拉回来再过滤。还有一点如果临时视图的表数据量很大尽量不要在外层对它做再次CTE或者多层嵌套因为Spark JDBC在生成最终SQL时会包一层子查询达梦对特别深的嵌套子查询支持得不是很好偶尔会报“临时表空间不足”或者“内存溢出”这类看着一头雾水的错。解决思路是让达梦算完该算的Spark只拿结果集做二次分析不要在Spark侧反复引用同一张达梦表做自连接。3.3 写回达梦追加与更新的正确姿势Spark把结果写回达梦最直观的方式就是resultDF.write .mode(append) .format(jdbc) .option(driver, dm.jdbc.driver.DmDriver) .option(url, jdbc:dm://127.0.0.1:5236) .option(user, SYSDBA) .option(password, ******) .option(dbtable, RESULT_TABLE) .option(batchsize, 2000) .save()mode支持append和overwrite。overwrite默认会先删表再重建表这里要非常小心Spark JDBC重建表时用的建表语句是通用SQL很多类型定义和达梦不完全匹配。如果目标表是已经存在的业务表我强烈建议不要用overwrite而是用append模式。如果确实需要先清理旧数据可以在Spark里先执行truncate table再走append这样比overwrite对表结构的破坏小得多。真正麻烦的是更新场景。Spark JDBC只支持表和行的batch写入不支持按主键update。我的做法是分两步第一步用Spark处理完结果写入一张临时表或者叫过渡表第二步用达梦的SQL把临时表的数据merge进业务主表。达梦的merge语法和Oracle很像MERGE INTO BUSINESS_TABLE t USING TEMP_TABLE s ON (t.id s.id) WHEN MATCHED THEN UPDATE SET t.name s.name, t.amount s.amount WHEN NOT MATCHED THEN INSERT (id, name, amount) VALUES (s.id, s.name, s.amount);如果数据量非常大可以对这个merge写成存储过程或批处理脚本分批次提交避免锁表时间太长影响在线业务。3.4 注册自定义方言根治SQL语法不兼容这是Spark适配达梦过程中最重要的点也是很多开发者在网上搜不到明确答案的地方。前面提到Spark对达梦的JDBC URL无法识别所以拿不到正确的方言导致生成SQL时可能用错标识符引用符、不支持某些函数甚至排序分页语法都出错。解决办法是手动注册一个达梦的方言类继承org.apache.spark.sql.jdbc.JdbcDialect重写两个方法import org.apache.spark.sql.jdbc.{JdbcDialect, JdbcDialects, JdbcType} import org.apache.spark.sql.types._ class DmJdbcDialect extends JdbcDialect { override def canHandle(url: String): Boolean url.startsWith(jdbc:dm:) override def quoteIdentifier(colName: String): String { // 达梦在默认模式下可以用双引号引用标识符 s${colName.replace(\, \\)} } override def getJDBCType(dt: DataType): Option[JdbcType] dt match { case StringType Some(JdbcType(VARCHAR2(4000), java.sql.Types.VARCHAR)) case BooleanType Some(JdbcType(NUMBER(1), java.sql.Types.NUMERIC)) case _ None } } // 在初始化代码里注册 JdbcDialects.registerDialect(new DmJdbcDialect())这段代码的核心逻辑是让Spark在拿到jdbc:dm://开头的URL时认为这是一个新的数据库类型并且用双引号来引用表名、字段名。达梦默认就使用双引号圈定标识符这和Oracle行为一致。同时通过重写getJDBCType可以让Spark在自动建表时把字符串映射成VARCHAR2而不是通用VARCHAR减少某些环境下的兼容问题。注册方言时要特别注意执行时机。如果在spark-shell里直接在启动后执行即可如果是Spark作业建议在main函数最开始执行保证后续所有读写操作都能识别到。如果作业里有多个入口比如用了spark-submit的多个类每个入口都要做一次注册。4. 性能调优与踩坑排查4.1 联调期最常遇到的连接问题我把实际联调过程中踩过的坑整理成了一张速查表基本覆盖了大多数人会遇到的问题现象可能原因解决方法ClassNotFoundException: dm.jdbc.driver.DmDriver驱动包未分发到Executor节点把jar放到Spark的jars目录或用--executor-extra-class-path指定Connection refused/ 连接超时达梦端口未开放、防火墙拦截检查5236端口连通性telnet IP 5236验证URL校验失败使用了jdbc:dm://以外的前缀确认URL格式为jdbc:dm://IP:端口连接上了但执行SQL报无效的列名字段名包含数据库关键字用自定义方言注册双引号引用或给列名加别名读取大表时Driver内存溢出fetchsize未设置默认全量拉取设置fetchsize并配合分区读取这里面最隐蔽的一个问题是在多节点环境下驱动jar只放在了提交作业的客户端节点结果Driver能正常连接但Executor一启动就报ClassNotFoundException。排查思路很简单找一个YARN节点的临时目录看看有没有这个jar或者直接在代码里打印一段Class.forName(dm.jdbc.driver.DmDriver)哪个节点报错就去查哪个节点。4.2 数据类型与关键字兼容的坑达梦和Spark之间最典型的类型兼容问题集中在三块大字段、时间类型、关键字。大字段就是CLOB和BLOB。Spark JDBC读取CLOB时如果列里存的是很长的文本可能出现读取不到内容或者乱码的问题。我一般是在达梦侧先把CLOB转成VARCHAR或者用DBMS_LOB.SUBSTR截断但要注意SUBSTR对超长文本会有长度限制。更好的办法是如果业务允许在Spark侧把该列转成StringType并指定大的字段长度然后通过自定义方言映射到VARCHAR2(4000)甚至TEXT达梦支持TEXT类型吗是的达梦兼容的TEXT类型也是存在的但我一般更倾向于用CLOB的底层存储映射。时间类型的坑主要在时区。Spark默认读取TIMESTAMP时会按JVM时区解析如果达梦服务器时区、Spark节点时区不一致拿到的时间会出现偏移。最好在连接URL或会话参数里统一时区并把Spark的spark.sql.session.timeZone设成和达梦一致。至于关键字达梦和Oracle一样对LEVEL、NUMBER、COMMENT、DATE这类词非常敏感。如果你的表字段名里恰好有这些词Spark生成的SQL里如果没有加双引号就会报SQL syntax error或者无效的标识符。这问题我刚接项目时几乎每天都在踩后来注册了自定义方言并在quoteIdentifier里统一处理才算彻底解决。4.3 分布式读写性能与稳定性调优最后聊一下把任务真正跑到集群上的性能表现。刚开始我做全量数据同步时发现作业跑了将近40分钟还没结束去达梦侧看数据库的AWR报告里出现了好几个长时间运行的SQL全是Spark生成的单条大查询。问题出在分区键的选择上——我用了ROW_NUMBER()生成的行号做分区键但这列上没有索引每个子查询都要排序后全表扫描代价非常高。后面我把分区键改成达梦表的主键并确保主键索引存在同时把numPartitions控制在10到20之间配合fetchsize1000整体同步时间从40分钟缩短到12分钟左右。这个提升巨大也说明了分区键的选择比分区数量更重要。写入侧的性能瓶颈主要在两个地方一是达梦的批次提交能力二是Spark写JDBC源时的并行度。batchsize不要设得过大我曾经设成10000结果达梦出现锁等待因为一个大事务长时间占用表级锁。建议1000到3000之间根据实际压测来调整。另外如果目标表上有多个索引大批量写入时索引维护成本很高可以考虑先drop索引写完再重新创建但这要结合业务的影响窗口来权衡。稳定性方面还有一个常见的报错是“未设置会话超时时间”。这个提示听起来像达梦参数问题实际上往往是达梦服务端在长时间空闲后主动断开了连接而Spark侧不知道等到下一次执行SQL时才报出来。解决方式是在达梦的配置文件里设置合适的会话超时参数或者确保Spark任务每次读写前都重新获取连接不要复用久置的旧连接。注意生产环境里对达梦实例做任何参数变更前一定要先和DBA确认避免影响其他在线业务。我一般会在测试库先验证再走变更流程推到生产。5. 一点后续可扩展的想法整个Spark适配达梦的实践到这里基本完整了。但落到工程上还有一个建议值得提不要把适配逻辑散落在各个业务代码里。我把达梦驱动、自定义方言、常用读写封装成了一个独立的模块团队的同学只要调用统一封装的方法传入目标表名和分区键就能完成读写不需要再关心底层方言这些问题。另外如果需要做增量同步可以考虑在达梦侧建一张增量日志表由Spark定期读取上次同步的游标位点来拉取增量数据。这个方案比每次全量读表高效得多也能大幅降低对达梦库的压力。我自己做这个项目时最大的体会是适配一种新的数据库最难的不是把数据拉通而是搞清楚连接器背后的各种隐藏规则。达梦作为国产数据库中应用很广的产品整体兼容度已经做得不错但Spark生态里它仍然是“小众玩家”所以需要我们这些使用者把链路补完整。希望这篇整理能帮到你如果你在适配过程中遇到这里没写到的问题欢迎私下多交流思路一起把这套方案做得更扎实。
返回列表