
1. 先搞清楚一个问题Spark访问达梦到底卡在哪大概两年前我接手了一个数据抽取需求业务库是达梦8数据量级在亿级数仓那边想用Spark做离线清洗和特征加工。当时第一反应很简单——Spark自带的JDBC数据源不是支持所有标准JDBC的数据库吗达梦作为国产数据库提供一个JDBC驱动不就完事了结果真正动手之后才发现事情远没有这么简单。很多做大数据开发的朋友第一次接触达梦时都会踩进同一个坑拿Spark连Oracle、MySQL的惯性思维去套达梦。实际上问题通常出在三个层面。第一层是驱动匹配问题达梦JDBC驱动和Spark内置的驱动加载逻辑之间经常出现类冲突或找不到驱动类第二层是SQL方言问题Spark JDBC数据源在读取数据时会自动生成一堆带双引号的SQL语句而达梦对大小写、引号、保留字的处理规则与Oracle和MySQL都不一样直接执行就报“表或视图不存在”第三层是性能和并发问题默认用单分区读取几亿行数据跑几个小时都不出结果而配置分区抽取时又会因为达梦连接数限制、事务机制差异导致连接风暴或写入失败。这篇文章就是把我在实际项目里把Spark和达梦“撮合”到一起的经验整理出来从驱动安装到读写调优再到排错思路完整过一遍。适合正准备做国产数据库适配的数据工程师、数仓开发以及被领导突然安排“把Spark任务迁移到达梦”的同学参考。先说结论Spark适配达梦这件事本身不难难在细节。达梦不是不能用Spark跑关键要对齐JDBC参数、注意方言规则、控制并发策略。下面按我实际操作的顺序展开。1.1 适配前的架构判断动手写代码之前我习惯先想清楚整个适配的架构Spark任务连达梦属于典型的外部数据源集成不走Spark DataSource API V2的定制插件完全依赖JDBC标准接口。所以最终的技术路径就是每个Spark Executor上的任务通过JDBC连接达梦实例Spark将SQL查询下推到达梦执行由达梦完成过滤和聚合后返回结果集再在Spark侧做分布式计算。这套方案的好处是改动最小不需要额外部署同步工具也不用把达梦数据先导出成文件再加载到HDFS。坏处是瓶颈一定在数据库侧的连接数和单条SQL执行效率上。理解了这个架构后面所有调优和排错就都有了方向——你在Spark端调多少并行度达梦那边就得撑住多少并发连接。1.2 可行的替换方案对比有朋友会问为什么不用DataX或者Sqoop做数据同步非要让Spark直连达梦我把当时评估过的几个方案放在一起对比过方案适合场景痛点DataX离线同步定时全量/增量搬运到HDFS或数仓额外维护同步任务实时性差做不到Spark SQL直接查达梦Sqoop导入Hadoop生态内的批量导入达梦不在官方支持列表需自己改插件维护成本高Spark JDBC直连需要Spark SQL直接关联查询、即席计算需要做好驱动、方言、并发适配先把数据导出成文件再处理小数据量、一次性任务亿级数据导出再导入周期太慢我最终选了Spark JDBC直连因为数仓那边还要用Spark SQL把达梦的维表和Hive里的行为表做关联计算这个场景只有直连能解决。如果你的需求只是单纯把达梦表同步到数仓DataX可能是更稳的选择。工具没有优劣匹配场景才是关键。2. 环境与依赖准备驱动放不对后面全是泪和很多“代码写半天不如环境配置折腾一天”的项目一样Spark适配达梦的第一步不是写代码而是把依赖环境收拾干净。这一节的内容全部来自我的实操经历每一步都踩过坑。2.1 达梦JDBC驱动的获取与确认先明确一点达梦官方提供的JDBC驱动和常见的MySQL、Oracle驱动在使用方式上差异不大但在版本兼容性上需要特别小心。我使用的是达梦8对应的高版本JDK环境驱动包是DmJdbcDriver18.jar在达梦安装目录的/dmdbms/drivers/jdbc下可以找到。拿到驱动包后不要急着丢进Spark里先用一个最基础的Java或命令行JDBC测试连一下确认驱动本身可用。这一步很多人会跳过导致后面出现问题时分不清到底是驱动问题还是Spark配置问题。# 在服务器上放置驱动并测试连接 mkdir -p /opt/dm_driver cp /dmdbms/drivers/jdbc/DmJdbcDriver18.jar /opt/dm_driver/ # 用一个最简单的JDBC测试程序验证连通性 java -cp /opt/dm_driver/DmJdbcDriver18.jar -Djava.security.egdfile:/dev/urandom \ -Djdbc.driversdm.jdbc.driver.DmDriver \ -jar dm_jdbc_test.jar测试通过后再考虑放进Spark环境。这里有一个容易被忽略的点驱动包不要随意改名称也不要和Spark自带的其他JDBC驱动混放在同一个自定义目录下避免出现类加载冲突。2.2 驱动放到哪个位置local模式与yarn集群模式的区别Spark读取JDBC数据源时驱动需要在Driver端和Executor端都可访问。最简单的做法是把驱动jar放到$SPARK_HOME/jars目录下这样Spark启动时会自动加载到classpath里。如果任务走的是yarn集群模式驱动jar的放置位置就有讲究了。仅放在提交机器的$SPARK_HOME/jars只能保证Driver端有驱动Executor在NodeManager上执行时照样报ClassNotFoundException。正确做法是提交任务时用--jars参数显式指定让yarn把jar分发到各个节点spark-submit \ --master yarn \ --deploy-mode cluster \ --jars /opt/dm_driver/DmJdbcDriver18.jar \ --driver-class-path /opt/dm_driver/DmJdbcDriver18.jar \ --conf spark.executor.extraClassPath/opt/dm_driver/DmJdbcDriver18.jar \ --class com.example.DmSparkApp \ spark-dm-adaptor.jar这三个参数每次提交任务都得检查一遍。我见过很多同事只加了--jars结果Driver端报No suitable driver原因就是--driver-class-path没带上。另外如果走的是Spark Thrift Server或者Livy驱动的放置规则又不同需要提前确认。2.3 用spark-sql快速验证连通性依赖放置好了我建议第一时间用spark-sql做一次最小化验证而不是直接写应用代码。最小化验证能把环境问题和业务代码问题隔离开-- 在spark-sql中创建临时视图指向达梦表 CREATE TEMPORARY VIEW dm_test USING jdbc OPTIONS ( url jdbc:dm://192.168.10.20:5236, dbtable TEST_USER.DIM_ORG, user test_user, password test_pass, driver dm.jdbc.driver.DmDriver ); -- 最简单的查询 SELECT * FROM dm_test LIMIT 10;如果这一步能出结果说明驱动加载、网络连通、JDBC URL、表定位都没问题接下来再考虑性能和并发。如果出不来结果就先排查网络端口、账号权限、驱动放置这三个层面。达梦默认端口是5236如果连不上可以先telnet 192.168.10.20 5236确认端口通不通排除防火墙干扰。3. 读取数据从打通连接到达梦数据的并行抽取连通性验证通过后真正的挑战才开始。Spark读取达梦并不是简单调一个load()就行分区策略、下推行为、类型映射都会直接影响任务能不能在可接受的时间内跑完。3.1 最小可用的读取示例与参数含义先给一个标准的Spark Scala读取示例这是最常用的写法val df spark.read .format(jdbc) .option(url, jdbc:dm://192.168.10.20:5236) .option(dbtable, TEST_USER.DIM_ORG) .option(user, test_user) .option(password, test_pass) .option(driver, dm.jdbc.driver.DmDriver) .option(fetchsize, 5000) .option(partitionColumn, org_id) .option(lowerBound, 1) .option(upperBound, 10000000) .option(numPartitions, 8) .load()这里每个参数都不是随便写的。fetchsize控制每次从达梦读取的行数默认值可能导致网络往返次数太多对亿级表来说影响明显partitionColumn、lowerBound、upperBound、numPartitions这四个参数共同决定Spark将查询拆分成多少个并行查询拆分列必须是有序的数值类型最常见的是主键。我遇到过有人拿字符串类型字段做分区列运行直接报错就是因为达梦的JDBC不支持下推这种分区谓词。3.2 分区列选不好并行就变成了灾难分区列的选择是读取性能的核心。理想情况下分区列应该满足三个条件数值类型、取值均匀、能建立索引。举个反面例子一张用户订单表有1亿条记录主键是order_id自增完全符合条件用分区读取可以同时开8个JDBC连接分别拉取不同区间的数据整个抽取时间可以缩小到串行方式的十分之一左右。但如果主键是varchar类型的UUID就没法直接做分区。这时候不能用UUID做partitionColumn常见的做法是用一个子查询作为dbtable在子查询里生成一个递增的序号列或者干脆以入库时间字段转换后的数值作为分区依据val dbtable (SELECT t.*, ROWNUM AS rn FROM TEST_USER.USER_ORDER t ) tmp 这种写法能绕开分区列必须是表中真实数值列的限制但代价是达梦要额外执行一次ROWNUM扫描需要评估子查询性能。如果原表本身有日期字段并且有索引直接在子查询里做日期范围过滤后再分区会更快。3.3 谓词下推别让Spark把整表数据都拉回来这是很多Spark新手最容易犯的错。默认情况下Spark JDBC数据源会尝试将查询条件下推到达梦执行但条件范围有限。比如你在Spark SQL里写了WHERE status validSpark只会下推简单的过滤条件而JOIN、GROUP BY这些操作大概率还是会先把数据拉到Spark内存里再处理。判断一个过滤条件是否真的下推成功可以直接打开Spark日志查看它打印出来的JDBC查询SQL。日志中如果显示SELECT ... WHERE status valid说明条件下推成功如果显示SELECT 全部字段 FROM table说明下推失败条件在Spark侧过滤。我在实际项目里发现达梦对Spark生成的下推SQL支持度还算可以但有个别情况需要注意如果过滤条件里涉及达梦特有的函数比如DECODESpark无法识别就不会下推。这种场景我的做法是在dbtable里用达梦原生SQL写一个子查询手动把过滤和裁剪做完.option(dbtable, (SELECT org_id, org_name, level_no FROM TEST_USER.DIM_ORG WHERE status valid AND level_no 2) tmp )手动子查询方式的可控性比依赖Spark自动下推更强尤其当表字段多、数据量大时提前裁剪字段可以显著降低网络传输量。3.4 类型映射带来的精度问题读取阶段还会遇到一个隐蔽的问题——类型映射。达梦的DECIMAL(38, 10)类型默认映射到Spark的DecimalType(38, 10)这是没有问题的。但有些老库表里的DECIMAL字段没有显式声明精度达梦默认给出的精度可能是DECIMAL(38, 10)和Spark的默认精度不一致就会导致结果集里出现精度丢失或转换异常。时间类型也一样。达梦的TIMESTAMP映射成Spark的TimestampType看起来没问题但如果达梦的会话时区与Spark的spark.sql.session.timeZone不一致读出来的时间会偏移。我的解决办法是在读取时使用dbtable子查询统一转成字符串或统一格式(SELECT TO_CHAR(create_time, YYYY-MM-DD HH24:MI:SS) AS create_time_str FROM TEST_USER.USER_ORDER) tmp虽然字符串后续处理要多一步转换但彻底避免了时区和格式不统一的坑。对于要求数据准确性优先的场景这样做是值得的。4. 写入数据批量落库与任务安全的平衡读取只是第一步Spark任务最终往往要把计算结果写回达梦供业务系统查询。写入比读取更容易出问题尤其是大批量写入时的事务、幂等性、批次大小这节我重点讲这些。4.1 SaveMode模式选择的实际差异Spark JDBC写入支持Append和Overwrite两种常用模式。很多人的第一反应是用Overwrite觉得省事但实际上在达梦上Overwrite的行为是Spark先执行DROP TABLE再创建新表如果表里有数据且有其他系统在引用这个操作轻则失败重则造成生产事故。我强烈建议生产环境一律使用Append模式表结构提前在达梦侧建好Spark只负责写数据df.write .mode(append) .format(jdbc) .option(url, jdbc:dm://192.168.10.20:5236) .option(dbtable, TEST_USER.ADS_ORG_LEVEL) .option(user, test_user) .option(password, test_pass) .option(driver, dm.jdbc.driver.DmDriver) .option(batchsize, 2000) .option(truncate, true) .option(isolationLevel, READ_COMMITTED) .save()如果确实需要覆盖写入正确的做法是先手动TRUNCATE TABLE再走Append写入。truncate参数在有的时候可以配合Overwrite使用但达梦JDBC下对truncate的支持需要实测我不建议依赖它。4.2 batchsize不是越大越好batchsize控制每个JDBC批处理提交的行数。默认值是1000对于达梦来说这个值可以适当调大但不要盲目拉到十万级。我实测过一张宽表50个字段的写入性能batchsize从500调到2000写入耗时下降明显但调到10000之后达梦JDBC驱动会报内存溢出或连接中断因为一个批次的数据量太大会撑爆驱动内部的缓冲。最终我稳定在2000到3000之间写入速度和稳定性达到了平衡。// 写入建议配置示例 .option(batchsize, 3000) .option(numPartitions, 6)numPartitions控制写入的并行度。如果并行度太高达梦同时收到大量写入会话很容易触发锁等待或连接数超限并行度太低大结果集写入又慢。我一般是先设6到8个分区观察数据库侧负载再逐步调整。4.3 任务重复执行的幂等性保障离线任务经常会有重跑场景。Spark写入天然不具备幂等性——同一个任务跑两次数据就重复一次。这个问题在达梦上比在MySQL上更隐蔽因为达梦默认事务隔离级别对并发写入的支持行为和预期不同。我的做法是引入一个“业务主键去重”策略-- 在达梦中创建目标表时带上唯一约束 CREATE TABLE TEST_USER.ADS_USER_ORDER ( order_id VARCHAR(50) PRIMARY KEY, user_id VARCHAR(50), order_amount DECIMAL(18,2), stat_date VARCHAR(10) ); -- 清理重复数据在重跑前先按主键删除历史分区数据 DELETE FROM TEST_USER.ADS_USER_ORDER WHERE stat_date 2024-05-20;每次任务启动时先按业务日期删除目标数据再做Append写入这样即使任务重跑多次最终结果也是一份干净的数据。比在Spark里用dropDuplicates后再写要高效得多因为Spark侧去重需要全量shuffle而数据库侧按日期删除用的是索引代价小很多。4.4 写入性能实测参考我在测试环境做了一组对比实验数据量是200万行、30个字段目标表在达梦侧已建好索引写入Spark的Executor配置为4个单executor 2核4G写入模式并行分区数batchsize耗时秒备注Serial单线程11000480极慢不推荐并行分区6100095正常水平并行分区6300068推荐配置并行分区12300051数据库CPU接近80%需谨慎并行分区610000失败驱动内存溢出能看到写入性能提升不是线性增长的并行度过高会让达梦成为瓶颈。我后来稳定用6到8个分区配合3000的batchsize写200万行大约1分钟出头能接受。5. 实录方言冲突、类型映射与分区膨胀的排错链路这一节把我在适配过程中遇到的几个典型问题完整复盘一遍。写出来不是为了展示过程而是想让你知道当报错信息出现时排查方向应该怎么定。5.1 “表或视图不存在”的真相大小写与双引号第一个遇到的坑是达梦里明明有TEST_USER.DIM_ORG这张表spark-sql读取时却报“表或视图不存在”。一开始我还以为是账号权限问题反复检查用户授权都没发现问题。后来我把Spark日志里生成的SQL打出来一看发现它生成的是SELECT ORG_ID, ORG_NAME FROM TEST_USER.DIM_ORG问题就出在双引号上。达梦的元数据默认将不带引号的表名、字段名转为大写存储但如果建表时带上了双引号就会保留原始大小写。Spark JDBC生成SQL时默认给标识符加双引号而达梦在双引号模式下会严格区分大小写。如果原表是用不带引号的CREATE TABLE语句创建的存的是大写名查询时用dim_org小写自然找不到。解决办法有几个一是在JDBC URL中加上compatibleModeoracle之类的参数部分版本能缓解大小写问题二是表名、字段名全部显式写成大写三是最稳的做法在dbtable选项里用子查询并加上小写别名绕过.option(dbtable, (SELECT org_id, org_name, create_time FROM TEST_USER.DIM_ORG) tmp)这样Spark生成的SQL会变成SELECT * FROM (SELECT ...) tmp不会对字段名加双引号也就绕开了大小写判断。这个技巧我一直沿用到项目结束基本没再犯过“表不存在”的错。5.2 驱动加载失败的完整排查链路另一个高频问题提交Spark任务后一直报java.sql.SQLException: No suitable driver found for jdbc:dm://...。这个报错看起来是驱动问题但根因可能是以下几种情况驱动jar没有被打包进Spark的classpath提交任务时只指定了--jars但yarn集群模式下Executor端没拿到jarSpark内置的驱动注册机制没有通过DriverManager注册达梦驱动类。排查思路我是这样做的先在提交任务的机器上手动执行一个Java进程确认JDBC连接串和驱动类名没问题然后把驱动jar同时放到$SPARK_HOME/jars和用--jars指定最后在代码里显式注册驱动Class.forName(dm.jdbc.driver.DmDriver)Class.forName这一行能解决90%的“No suitable driver”问题。因为Spark的驱动管理器在部分部署方式下不会自动加载第三方驱动的META-INF服务描述文件显式注册最省心。5.3 DECIMAL精度和TIMESTAMP时区偏移的记录类型问题更多出现在没有统一规范的老表上。我们有一张表的金额字段在达梦里的类型是DECIMAL没带精度默认成了DECIMAL(38, 10)Spark读出来后在结果表里写回整数金额时出现了小数点错位。定位方式是比较达梦侧的原始数据和Spark算出来的结果发现每个数字都差了10的次方倍。根因是Spark从JDBC结果集里获取DECIMAL精度时依赖的是数据库元数据达梦对未声明精度的字段返回的精度值在驱动实现上和Spark预期不一致。解决办法就是在dbtable子查询里用CAST显式指定(SELECT CAST(order_amount AS DECIMAL(18,2)) AS order_amount FROM TEST_USER.USER_ORDER) tmp时间字段的偏移也是类似处理。统一在达梦侧用TO_CHAR或TO_DATE转成目标格式Spark侧只当字符串处理绕开时区转换。5.4 distinct之后分区数暴涨的意外还有一个比较隐蔽的问题我用spark.sql(SELECT DISTINCT user_id FROM dm_table)对达梦读取出来的DataFrame做去重处理完后DataFrame分区数从8个暴涨到了200多个导致后续写入达梦时并行写入压力骤增数据库侧的锁竞争非常严重。原因在于DISTINCT操作在Spark底层会引入一个hashAggregate为了并行计算会重新分区。这不是达梦特有但很多人在做数据库适配时容易忽略这层。解决办法很简单在DISTINCT之后显式coalesce或repartition回合理的分区数val result df.select(user_id).distinct().coalesce(6)或者更优雅的做法是直接在dbtable里让达梦做完去重Spark只负责拉取结果.option(dbtable, (SELECT DISTINCT user_id FROM TEST_USER.USER_ORDER) tmp)这样既减少Spark shuffle也减少写入并发。6. 资源参数与运行期监控的实测心得最后聊一聊上线之后的运行参数调整和监控手段。适配做完只是第一步跑得稳不稳、达梦会不会被Spark任务拖垮才是生产环境真正考验的地方。6.1 连接数与连接风暴的克制Spark JDBC读取时每个分区在Executor里都会建立独立的JDBC连接。假设某个任务有20个分区分布在10个Executor上达梦同一时间就会收到20个连接。如果同时有多个Spark任务在跑连接数很容易冲破达梦实例的最大连接限制。达梦默认最大连接数不一定能满足生产需要需要在dm.ini里调整MAX_SESSIONS参数。但与其无限调大数据库连接数更推荐在Spark侧做资源管控。我给团队定的规范是场景并行分区数同时运行任务数说明小时级ETL81~2控制并发避免锁表日级全量抽取121大表抽取时独占连接资源临时查询42~3控制即时查询对库的影响6.2 从Spark UI和达梦动态视图双端观察运行期我一般两头看一头是Spark UI的Task耗时、Shuffle大小、GC时间另一头是达梦侧的会话情况。曾有次写入任务跑得很慢Spark UI显示每个Task都在疯狂GC通过达梦动态视图一查发现数据库端有大量会话处于ACTIVE状态等待锁释放说明任务卡在了数据库锁上。在达梦上常用的查询语句-- 查看当前会话及状态 SELECT * FROM V$SESSIONS WHERE STATE ACTIVE; -- 查看锁等待情况 SELECT * FROM V$LOCK WHERE BLOCKED 1;如果发现大量锁等待优先降低Spark写入并行度其次检查代码里是否有跨任务共享事务的情况。我在一次大批量写入时就是靠这个查询锁定了问题——两个Spark任务同时写同一张表互相阻塞调整调度时间后马上恢复。6.3 连接空闲超时与长任务稳定性达梦的数据库会话默认有超时时间但很多环境的参数并没有显式配置。如果一个Spark任务从读取到计算再到写入单条JDBC连接的空闲时间超过数据库阈值连接可能被数据库端强制断开任务随之报错。这种情况我在抽取大表时碰到过一次。解决方案是在JDBC URL里显式设置连接超时和会话参数jdbc:dm://192.168.10.20:5236?socketTimeout300000connectionTimeout60000同时也建议在达梦侧确认TCP_CONN_TIMEOUT、LOGIN_TIMEOUT等参数不是过小值。两边设置配平后长任务稳定多了。最后分享一个我自己的习惯每次调整完Spark访问达梦的参数我都会在测试环境用同一张500万行的表跑一遍基线测试记录读取耗时、写入耗时、达梦侧CPU峰值。这样任何一次参数改动是变好还是变差一目了然不会靠感觉调参。适配国产数据库这件事大多时候不是技术门槛高而是细节太多——驱动版本、大小写规则、连接数上限、事务行为每一项都值得较真。希望这篇分享能让你少走几趟弯路。