ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

达梦数据库存储过程与定时任务实现数据自动迁移方案

达梦数据库存储过程与定时任务实现数据自动迁移方案 数据迁移这件事做过一次的人都知道最怕的不是迁移本身而是迁移完之后业务方隔三差五来找你“昨天的数据怎么还没同步过来”你打开工具手动跑一遍数据是过来了但明天呢后天呢这种重复劳动做多了人就会开始琢磨——能不能让数据库自己把这件事干了达梦数据库作为国产数据库中的主力选手在越来越多的生产环境中承担着核心业务数据的存储和管理。但很多团队在从其他数据库迁移到达梦之后往往只关注“数据能不能导进去”而忽略了“数据能不能持续自动地流转”。手动导数在一次性迁移场景下没问题可一旦涉及周期性同步、跨库数据汇总、历史数据归档这类需求靠人工操作就完全不现实了。这篇文章要聊的就是怎么利用达梦的存储过程配合定时任务搭建一套能自己跑起来的数据自动迁移方案。这套方案适合谁如果你正在使用达梦数据库手头有周期性数据同步的需求又不想引入额外的ETL工具或中间件那这篇文章的内容可以直接拿来参考。如果你对存储过程和定时任务还不算熟悉也没关系我会把每一步的逻辑和踩过的坑都讲清楚。1. 为什么选存储过程加定时任务这条路线1.1 达梦环境下数据迁移的几种常见做法在达梦数据库里做数据迁移摆在面前的路其实有好几条。最直接的是用达梦自带的数据迁移工具DTS图形化界面点几下就能把源库的表结构和数据搬到目标库。这种方式适合一次性迁移比如系统上线前的数据割接。但它的局限也很明显——你没法让它每天凌晨自动跑一次除非你每天凌晨爬起来点鼠标。另一条路是写外部脚本比如用Python或者Shell脚本连接达梦执行查询和插入操作再用操作系统的crontab来调度。这种方式灵活度高但维护成本也高。脚本散落在各个服务器上时间一长谁写的、干什么用的、依赖什么环境全都说不清楚。而且脚本里的数据库连接信息一旦变更你得挨个去改。还有一条路就是本文要讲的把迁移逻辑封装在达梦的存储过程里再用数据库自身的定时任务机制来调度。这条路线的好处是迁移逻辑和调度逻辑都在数据库内部不依赖外部脚本和操作系统级的调度器。数据库能跑任务就能跑。备份恢复的时候存储过程和调度配置一起跟着走不会出现“数据恢复了但任务没恢复”的尴尬。1.2 存储过程封装迁移逻辑的天然优势存储过程这东西很多人觉得它老派、不好维护但在数据迁移这个场景下它有几个别人替代不了的优点。第一个是事务控制。数据迁移最怕的就是迁了一半失败了源数据删了目标数据没进去。存储过程里可以用事务把一批操作包起来要么全成功要么全回滚不会留下中间状态。这一点用外部脚本做也不是不行但脚本里的连接管理和事务边界控制写起来比存储过程麻烦得多。第二个是执行效率。存储过程在数据库内部执行省去了网络往返的开销。特别是涉及大批量数据的时候存储过程可以分批提交避免一次性锁太多行导致日志暴涨。你可以控制每批处理多少条批与批之间做一次提交既保证了效率又不至于把undo空间撑爆。第三个是权限收敛。你不需要把数据库的账号密码散落在各个脚本里只需要给调度任务配置一个执行存储过程的权限就行。存储过程内部的操作对调用者来说是透明的调用者只需要知道“执行这个存储过程”不需要知道它具体动了哪些表。1.3 DBMS_SCHEDULER在达梦中的角色定位达梦数据库提供了DBMS_SCHEDULER包这是从Oracle体系延续下来的任务调度机制。你可以把它理解成数据库内置的一个“闹钟服务”——你告诉它什么时候执行、执行什么、执行频率是多少它就按时去干活。和操作系统层面的crontab相比DBMS_SCHEDULER的优势在于它和数据库是同一套体系。任务的执行记录、执行状态、错误信息都记录在数据库的视图里你可以直接查DBA_SCHEDULER_JOB_LOG来看某个任务昨天有没有跑成功。而crontab的日志散落在系统日志里查起来没那么方便。另外DBMS_SCHEDULER支持更细粒度的调度策略。比如你可以设置任务在工作日的每小时执行一次但避开整点的高峰期也可以设置任务在某个特定日期之后才开始生效。这些在crontab里虽然也能实现但配置起来要绕一些。注意达梦的DBMS_SCHEDULER在语法上和Oracle高度兼容但并非100%一致。在实际使用前建议先查一下当前达梦版本的官方文档确认支持的参数范围。2. 迁移存储过程的设计与编写要点2.1 先想清楚迁移的粒度全量还是增量写存储过程之前第一个要回答的问题是每次执行的时候迁移哪些数据如果是全量迁移逻辑很简单——清空目标表把源表所有数据插进去。但这种做法在数据量大的时候非常危险一次迁移可能跑几个小时期间目标表处于不可用状态。而且如果迁移过程中出现网络抖动或者锁等待整个操作回滚前功尽弃。更常见的做法是增量迁移。每次只迁移上次迁移之后新增或变更的数据。这就需要一个“水位线”的概念——记录上次迁移到了哪个时间点或者哪个ID。下次执行的时候从水位线之后开始取数。水位线的维护方式有几种。最简单的是在目标库建一张配置表每次迁移完成后更新水位线字段。另一种方式是利用源表本身的时间戳字段比如UPDATE_TIME每次取UPDATE_TIME 上次迁移时间的记录。后者的前提是源表有可靠的时间戳维护机制如果源表的UPDATE_TIME不可信那就只能用配置表的方式。-- 水位线配置表示例 CREATE TABLE MIGRATION_WATERMARK ( TASK_NAME VARCHAR(100) PRIMARY KEY, LAST_VALUE VARCHAR(50), UPDATE_TIME TIMESTAMP DEFAULT SYSDATE ); -- 初始化一条水位线记录 INSERT INTO MIGRATION_WATERMARK (TASK_NAME, LAST_VALUE) VALUES (ORDER_SYNC, 2024-01-01 00:00:00); COMMIT;2.2 分批提交避免大事务把日志撑爆增量迁移虽然每次的数据量比全量小但如果业务高峰期积累了几个小时的数据一次迁移的量也可能很大。这时候如果用一个事务把所有插入做完undo日志会急剧膨胀严重的时候可能把表空间撑满。解决办法是分批提交。在存储过程里用循环每次取一批数据比如1000条插入目标表后提交一次然后继续取下一批。这样即使中途失败已经提交的批次不会丢失下次执行的时候从水位线继续就行。CREATE OR REPLACE PROCEDURE PROC_MIGRATE_ORDER( P_BATCH_SIZE IN INT DEFAULT 1000 ) AS V_LAST_TIME TIMESTAMP; V_MAX_TIME TIMESTAMP; V_ROW_COUNT INT; BEGIN -- 获取当前水位线 SELECT LAST_VALUE INTO V_LAST_TIME FROM MIGRATION_WATERMARK WHERE TASK_NAME ORDER_SYNC; -- 循环分批迁移 LOOP -- 取源表当前批次的最大时间 SELECT MAX(UPDATE_TIME) INTO V_MAX_TIME FROM ( SELECT UPDATE_TIME FROM SRC_ORDER WHERE UPDATE_TIME V_LAST_TIME ORDER BY UPDATE_TIME LIMIT P_BATCH_SIZE ); EXIT WHEN V_MAX_TIME IS NULL; -- 插入目标表 INSERT INTO TGT_ORDER (ORDER_ID, ORDER_NO, AMOUNT, UPDATE_TIME) SELECT ORDER_ID, ORDER_NO, AMOUNT, UPDATE_TIME FROM SRC_ORDER WHERE UPDATE_TIME V_LAST_TIME AND UPDATE_TIME V_MAX_TIME; V_ROW_COUNT : SQL%ROWCOUNT; COMMIT; -- 更新水位线 UPDATE MIGRATION_WATERMARK SET LAST_VALUE TO_CHAR(V_MAX_TIME, YYYY-MM-DD HH24:MI:SS), UPDATE_TIME SYSDATE WHERE TASK_NAME ORDER_SYNC; COMMIT; V_LAST_TIME : V_MAX_TIME; EXIT WHEN V_ROW_COUNT P_BATCH_SIZE; END LOOP; EXCEPTION WHEN OTHERS THEN ROLLBACK; -- 记录错误日志 INSERT INTO MIGRATION_LOG (TASK_NAME, ERR_CODE, ERR_MSG, LOG_TIME) VALUES (ORDER_SYNC, SQLCODE, SQLERRM, SYSDATE); COMMIT; RAISE; END; /这段代码有几个细节值得展开说。关于LIMIT的使用达梦支持LIMIT语法但在子查询里配合ORDER BY使用的时候要注意如果源表数据量很大每次取MAX(UPDATE_TIME)都要扫一遍排序效率会随着数据量增长而下降。优化思路是在UPDATE_TIME上建索引让排序走索引扫描。关于水位线更新和业务插入的提交顺序上面的代码是先插入业务数据、提交再更新水位线、提交。这两个提交之间如果发生故障会出现“数据已迁移但水位线没更新”的情况下次执行会重复迁移一批数据。解决办法是在目标表上建唯一约束重复插入的时候用MERGE或者INSERT ... ON DUPLICATE KEY来处理。或者把水位线更新和业务插入放在同一个事务里但这样又回到了大事务的问题。实际项目中通常选择“允许少量重复用唯一约束兜底”的方案。关于异常处理存储过程里的EXCEPTION块捕获异常后先把错误信息写入日志表并提交然后再RAISE把异常抛出去。这样做的目的是让调度任务能感知到执行失败同时错误信息不会因为回滚而丢失。2.3 字段映射和类型转换的坑跨库迁移的时候源表和目标表的字段类型往往不完全一致。比如源库的DATE类型到了达梦可能对应TIMESTAMP源库的VARCHAR2到了达梦可能对应VARCHAR。这些类型差异在简单查询的时候可能看不出来但在存储过程里做插入的时候如果类型不匹配轻则隐式转换导致精度丢失重则直接报错。达梦在类型转换上比Oracle要严格一些。比如把字符串2024-01-01插入DATE类型的字段Oracle可能会自动转换但达梦在某些版本下会报错。稳妥的做法是在存储过程里显式转换-- 显式转换避免隐式转换的坑 INSERT INTO TGT_ORDER (ORDER_ID, ORDER_DATE, AMOUNT) SELECT ORDER_ID, TO_DATE(ORDER_DATE_STR, YYYY-MM-DD), CAST(AMOUNT AS DECIMAL(18,2)) FROM SRC_ORDER WHERE UPDATE_TIME V_LAST_TIME;还有一个容易忽略的点是空字符串和NULL的区别。在Oracle里空字符串和NULL是等价的但在达梦里某些版本下空字符串会被当作一个独立的空串值处理。如果源表里有空字符串迁移到达梦后可能会变成NULL导致目标表的非空约束报错。处理办法是在插入前用NVL或者CASE WHEN把空串转成默认值。2.4 迁移过程中的锁与性能平衡存储过程执行期间源表如果有大量的写入操作可能会出现锁等待。特别是当迁移的查询条件命中了源表的索引而业务写入也在操作同一批数据的时候锁冲突的概率会明显上升。一个实用的技巧是控制迁移的时间窗口。把定时任务安排在业务低峰期执行比如凌晨2点到4点之间。这样即使迁移过程对源表加了共享锁也不会对业务造成明显影响。另一个技巧是使用游标分批读取。上面的示例代码用的是LIMIT分批这种方式在数据量适中的时候没问题。但如果源表数据量特别大每次SELECT MAX(UPDATE_TIME)都要扫描大量数据效率会下降。这时候可以改用游标DECLARE CURSOR C_ORDER IS SELECT ORDER_ID, ORDER_NO, AMOUNT, UPDATE_TIME FROM SRC_ORDER WHERE UPDATE_TIME V_LAST_TIME ORDER BY UPDATE_TIME; V_ORDER C_ORDER%ROWTYPE; V_COUNT INT : 0; BEGIN OPEN C_ORDER; LOOP FETCH C_ORDER INTO V_ORDER; EXIT WHEN C_ORDER%NOTFOUND; INSERT INTO TGT_ORDER VALUES ( V_ORDER.ORDER_ID, V_ORDER.ORDER_NO, V_ORDER.AMOUNT, V_ORDER.UPDATE_TIME ); V_COUNT : V_COUNT 1; IF MOD(V_COUNT, 1000) 0 THEN COMMIT; END IF; END LOOP; CLOSE C_ORDER; COMMIT; END;游标方式的好处是只扫描一次源表后续的读取都在游标缓冲区里进行。但缺点是游标会持有源表的读一致性快照如果迁移时间很长undo表空间会持续增长。所以游标方式适合数据量中等、迁移窗口充裕的场景。3. 用DBMS_SCHEDULER把存储过程挂上定时器3.1 创建Job的基本语法和参数解读存储过程写好了接下来就是让它自动跑起来。达梦的DBMS_SCHEDULER创建Job的基本语法如下BEGIN DBMS_SCHEDULER.CREATE_JOB( JOB_NAME JOB_MIGRATE_ORDER, JOB_TYPE STORED_PROCEDURE, JOB_ACTION PROC_MIGRATE_ORDER, START_DATE SYSDATE, REPEAT_INTERVAL FREQDAILY; BYHOUR2; BYMINUTE0; BYSECOND0, ENABLED TRUE, COMMENTS 订单数据每日凌晨2点自动迁移 ); END; /几个关键参数需要解释一下。JOB_TYPE指定任务类型这里用的是STORED_PROCEDURE表示执行一个存储过程。达梦还支持PLSQL_BLOCK可以直接写一段PL/SQL代码块适合逻辑比较简单的场景。如果迁移逻辑不复杂用PLSQL_BLOCK可以省去创建存储过程的步骤。REPEAT_INTERVAL是调度频率用的是日历表达式语法。FREQDAILY; BYHOUR2; BYMINUTE0; BYSECOND0表示每天凌晨2点整执行。这个表达式可以组合出很复杂的调度策略比如FREQWEEKLY; BYDAYMON,TUE,WED,THU,FRI; BYHOUR3表示周一到周五每天凌晨3点执行。ENABLED参数设为TRUE表示创建后立即启用。如果设为FALSE任务创建后处于禁用状态需要手动调用DBMS_SCHEDULER.ENABLE来启用。建议在生产环境中先设为FALSE等确认配置无误后再启用。3.2 调度频率的实战配置别让任务撞上业务高峰调度频率的设置看起来简单但实际配置的时候有几个坑。第一个坑是任务执行时间超过了调度间隔。比如你设置每10分钟执行一次但某次执行因为数据量突增跑了15分钟。这时候下一次调度时间已经到了但上一次还没跑完。达梦默认的行为是等待上一次执行完成后再执行下一次但这会导致任务堆积。解决办法是在存储过程开头加一个“是否正在执行”的判断如果上一次还没跑完本次直接跳过。CREATE OR REPLACE PROCEDURE PROC_MIGRATE_ORDER(...) AS V_RUNNING INT; BEGIN -- 检查是否有正在执行的同类任务 SELECT COUNT(*) INTO V_RUNNING FROM MIGRATION_LOG WHERE TASK_NAME ORDER_SYNC AND STATUS RUNNING AND LOG_TIME SYSDATE - 1/24; -- 1小时内 IF V_RUNNING 0 THEN RETURN; -- 上一次还在跑本次跳过 END IF; -- 记录开始执行 INSERT INTO MIGRATION_LOG (TASK_NAME, STATUS, LOG_TIME) VALUES (ORDER_SYNC, RUNNING, SYSDATE); COMMIT; -- ... 迁移逻辑 ... -- 记录执行完成 UPDATE MIGRATION_LOG SET STATUS SUCCESS WHERE TASK_NAME ORDER_SYNC AND STATUS RUNNING; COMMIT; END; /第二个坑是调度表达式的时间基准。达梦的START_DATE和REPEAT_INTERVAL是配合使用的。如果START_DATE设的是SYSDATE而当前时间是下午3点REPEAT_INTERVAL设的是FREQDAILY; BYHOUR2那么第一次执行会在明天凌晨2点而不是今天。这个行为是符合预期的但如果不注意可能会误以为任务没生效。第三个坑是时区问题。如果数据库服务器和应用服务器不在同一个时区调度时间的计算可能会出现偏差。达梦的调度时间是基于数据库服务器的时间所以在配置之前先用SELECT SYSDATE FROM DUAL确认一下数据库的当前时间。3.3 任务执行状态的监控与日志查询任务创建之后怎么知道它有没有在跑、跑得怎么样达梦提供了一系列视图来查看调度任务的状态。视图名称用途DBA_SCHEDULER_JOBS查看所有调度任务的配置信息DBA_SCHEDULER_JOB_LOG查看任务执行的历史日志DBA_SCHEDULER_RUNNING_JOBS查看当前正在执行的任务DBA_SCHEDULER_JOB_RUN_DETAILS查看任务执行的详细结果查任务配置SELECT JOB_NAME, JOB_TYPE, JOB_ACTION, REPEAT_INTERVAL, ENABLED, STATE FROM DBA_SCHEDULER_JOBS WHERE JOB_NAME JOB_MIGRATE_ORDER;查最近10次执行记录SELECT LOG_DATE, STATUS, ACTUAL_START_DATE, RUN_DURATION, ADDITIONAL_INFO FROM DBA_SCHEDULER_JOB_LOG WHERE JOB_NAME JOB_MIGRATE_ORDER ORDER BY LOG_DATE DESC LIMIT 10;如果发现任务执行失败ADDITIONAL_INFO字段里会有错误信息。常见的失败原因包括存储过程不存在、权限不足、源表被锁等。根据错误信息定位问题后修复并重新启用任务即可。提示建议在存储过程内部也维护一张自定义的日志表记录每次迁移的开始时间、结束时间、迁移行数、错误信息等。这样即使调度视图的日志被清理了你仍然有完整的迁移历史可查。4. 那些只有踩过才知道的坑4.1 权限配置别让任务因为权限问题静默失败DBMS_SCHEDULER创建的任务默认以创建者的身份执行。如果创建Job的用户和存储过程的所有者不是同一个用户或者存储过程内部访问了其他Schema的表就可能出现权限不足的问题。最稳妥的做法是用存储过程的所有者来创建Job。如果做不到那就需要显式授予权限。比如存储过程内部要访问SRC_SCHEMA.SRC_ORDER表那么Job的创建者需要对这张表有SELECT权限。还有一个容易被忽略的点是CREATE JOB权限。普通用户默认没有创建调度任务的权限需要DBA授予GRANT CREATE JOB TO USER_MIGRATION;另外如果存储过程内部有INSERT、UPDATE、DELETE操作还需要确保Job执行者有对应的对象权限。权限问题最麻烦的地方在于它往往不会在创建Job的时候报错而是在任务实际执行的时候才失败。所以创建完Job之后建议手动执行一次DBMS_SCHEDULER.RUN_JOB来验证权限是否配置正确。4.2 迁移过程中的数据类型精度丢失前面提到了类型转换的问题这里再展开说一个更隐蔽的坑数值精度丢失。假设源库的金额字段是NUMBER(10,2)目标库的金额字段是DECIMAL(18,2)看起来目标库的精度更大应该没问题。但如果源库的金额字段实际存储了NUMBER(10,4)的数据迁移到达梦的DECIMAL(18,2)字段时小数部分会被截断。这种精度丢失在迁移过程中不会报错但迁移完成后对账的时候就会发现金额对不上。解决办法是在迁移前先做一次数据探查确认源表每个字段的实际数据精度然后确保目标表的字段精度不低于源表。如果目标表的精度确实需要缩小那就要在存储过程里显式做四舍五入而不是依赖数据库的隐式截断。-- 显式四舍五入避免隐式截断 INSERT INTO TGT_ORDER (ORDER_ID, AMOUNT) SELECT ORDER_ID, ROUND(AMOUNT, 2) FROM SRC_ORDER;4.3 任务执行时间过长导致的连锁反应一个迁移任务如果执行时间过长可能会引发一系列连锁问题。首先是锁等待。迁移任务对源表的查询会加共享锁如果业务此时要对同一批数据做更新就会出现锁等待。等待时间长了业务连接池可能被占满进而影响整个应用。其次是日志空间。迁移过程中的插入操作会产生redo日志如果迁移量很大redo日志文件可能会被快速写满触发日志切换。如果归档空间不足数据库会挂起。再次是调度堆积。如果任务执行时间超过了调度间隔下一次调度触发的时候上一次还没结束任务就会排队等待。排队多了之后调度器可能会报错。应对这些问题的策略是控制单次迁移的数据量。在存储过程里设置一个最大迁移行数比如每次最多迁移10万行超过的部分留到下一次执行。这样单次执行时间可控不会对数据库造成太大压力。-- 设置单次最大迁移行数 V_MAX_ROWS INT : 100000; V_TOTAL_ROWS INT : 0; LOOP -- ... 分批迁移 ... V_TOTAL_ROWS : V_TOTAL_ROWS V_ROW_COUNT; EXIT WHEN V_TOTAL_ROWS V_MAX_ROWS; END LOOP;4.4 达梦版本差异带来的兼容性问题达梦数据库的版本迭代比较快不同版本之间在DBMS_SCHEDULER和存储过程语法上可能存在差异。比如某些早期版本不支持LIMIT语法需要用ROWNUM来替代某些版本对REPEAT_INTERVAL的日历表达式支持不完整。在实际项目中我遇到过达梦7和达梦8在存储过程异常处理上的行为差异。达梦7里SQLCODE和SQLERRM在某些异常场景下返回的值和达梦8不一致导致错误日志记录不准确。解决办法是在存储过程里用WHEN OTHERS THEN捕获异常后同时记录SQLCODE和自定义的错误描述而不是完全依赖SQLERRM。另一个版本差异是调度任务的时区处理。达梦8的某些版本在START_DATE的处理上会考虑数据库的时区设置而达梦7则统一按服务器本地时间处理。如果从达梦7升级到达梦8原有的调度任务可能需要重新调整时间配置。建议在正式部署之前先在测试环境用目标版本做一次完整的验证包括存储过程编译、Job创建、手动执行、自动调度等环节。不要假设在开发环境能跑通的代码在生产环境也一定能跑通。5. 从手动到自动一个完整的落地案例5.1 场景描述与表结构设计假设有一个订单系统源库是MySQL目标库是达梦。业务要求每天凌晨把MySQL中当天的订单数据同步到达梦的分析库中供报表系统查询。源表结构MySQLCREATE TABLE src_order ( order_id BIGINT PRIMARY KEY, order_no VARCHAR(32), customer_id BIGINT, amount DECIMAL(12,2), order_status TINYINT, create_time DATETIME, update_time DATETIME );目标表结构达梦CREATE TABLE TGT_ORDER ( ORDER_ID BIGINT PRIMARY KEY, ORDER_NO VARCHAR(32), CUSTOMER_ID BIGINT, AMOUNT DECIMAL(12,2), ORDER_STATUS INT, CREATE_TIME TIMESTAMP, UPDATE_TIME TIMESTAMP, SYNC_TIME TIMESTAMP DEFAULT SYSDATE );注意目标表多了一个SYNC_TIME字段用来记录这条数据是什么时候同步过来的。这个字段在排查问题的时候非常有用——如果发现某天的数据有问题可以直接查SYNC_TIME来定位是哪次迁移带过来的。5.2 存储过程的完整实现CREATE OR REPLACE PROCEDURE PROC_SYNC_ORDER( P_BATCH_SIZE IN INT DEFAULT 2000, P_MAX_ROWS IN INT DEFAULT 100000 ) AS V_LAST_TIME TIMESTAMP; V_MAX_TIME TIMESTAMP; V_ROW_COUNT INT; V_TOTAL_ROWS INT : 0; V_START_TIME TIMESTAMP; BEGIN V_START_TIME : SYSDATE; -- 记录任务开始 INSERT INTO MIGRATION_LOG (TASK_NAME, STATUS, LOG_TIME) VALUES (ORDER_SYNC, RUNNING, V_START_TIME); COMMIT; -- 获取水位线 BEGIN SELECT LAST_VALUE INTO V_LAST_TIME FROM MIGRATION_WATERMARK WHERE TASK_NAME ORDER_SYNC; EXCEPTION WHEN NO_DATA_FOUND THEN V_LAST_TIME : TO_TIMESTAMP(2024-01-01 00:00:00, YYYY-MM-DD HH24:MI:SS); INSERT INTO MIGRATION_WATERMARK (TASK_NAME, LAST_VALUE) VALUES (ORDER_SYNC, TO_CHAR(V_LAST_TIME, YYYY-MM-DD HH24:MI:SS)); COMMIT; END; -- 分批迁移 LOOP EXIT WHEN V_TOTAL_ROWS P_MAX_ROWS; -- 获取当前批次的最大时间 SELECT MAX(UPDATE_TIME) INTO V_MAX_TIME FROM ( SELECT UPDATE_TIME FROM SRC_ORDER WHERE UPDATE_TIME V_LAST_TIME ORDER BY UPDATE_TIME LIMIT P_BATCH_SIZE ); EXIT WHEN V_MAX_TIME IS NULL; -- 插入目标表用MERGE避免重复 MERGE INTO TGT_ORDER T USING ( SELECT ORDER_ID, ORDER_NO, CUSTOMER_ID, AMOUNT, ORDER_STATUS, CREATE_TIME, UPDATE_TIME FROM SRC_ORDER WHERE UPDATE_TIME V_LAST_TIME AND UPDATE_TIME V_MAX_TIME ) S ON (T.ORDER_ID S.ORDER_ID) WHEN MATCHED THEN UPDATE SET T.ORDER_NO S.ORDER_NO, T.AMOUNT S.AMOUNT, T.ORDER_STATUS S.ORDER_STATUS, T.UPDATE_TIME S.UPDATE_TIME, T.SYNC_TIME SYSDATE WHEN NOT MATCHED THEN INSERT (ORDER_ID, ORDER_NO, CUSTOMER_ID, AMOUNT, ORDER_STATUS, CREATE_TIME, UPDATE_TIME, SYNC_TIME) VALUES (S.ORDER_ID, S.ORDER_NO, S.CUSTOMER_ID, S.AMOUNT, S.ORDER_STATUS, S.CREATE_TIME, S.UPDATE_TIME, SYSDATE); V_ROW_COUNT : SQL%ROWCOUNT; COMMIT; -- 更新水位线 UPDATE MIGRATION_WATERMARK SET LAST_VALUE TO_CHAR(V_MAX_TIME, YYYY-MM-DD HH24:MI:SS), UPDATE_TIME SYSDATE WHERE TASK_NAME ORDER_SYNC; COMMIT; V_LAST_TIME : V_MAX_TIME; V_TOTAL_ROWS : V_TOTAL_ROWS V_ROW_COUNT; EXIT WHEN V_ROW_COUNT P_BATCH_SIZE; END LOOP; -- 记录任务完成 UPDATE MIGRATION_LOG SET STATUS SUCCESS, ROW_COUNT V_TOTAL_ROWS, DURATION (SYSDATE - V_START_TIME) * 86400 WHERE TASK_NAME ORDER_SYNC AND STATUS RUNNING AND LOG_TIME V_START_TIME; COMMIT; EXCEPTION WHEN OTHERS THEN ROLLBACK; INSERT INTO MIGRATION_LOG (TASK_NAME, STATUS, ERR_CODE, ERR_MSG, LOG_TIME) VALUES (ORDER_SYNC, FAILED, SQLCODE, SQLERRM, SYSDATE); COMMIT; RAISE; END; /这段代码比前面的示例更完整加入了MERGE语句来处理重复数据加入了任务开始和结束的日志记录加入了执行时长的统计。MERGE语句是达梦支持的它的好处是“存在则更新不存在则插入”避免了先删后插或者先查后插的繁琐逻辑。5.3 创建调度任务并验证存储过程编译通过后创建调度任务BEGIN DBMS_SCHEDULER.CREATE_JOB( JOB_NAME JOB_SYNC_ORDER, JOB_TYPE STORED_PROCEDURE, JOB_ACTION PROC_SYNC_ORDER, START_DATE TRUNC(SYSDATE) 1 2/24, -- 明天凌晨2点 REPEAT_INTERVAL FREQDAILY; BYHOUR2; BYMINUTE0; BYSECOND0, ENABLED FALSE, COMMENTS 订单数据每日凌晨2点同步 ); END; /先不启用手动执行一次验证BEGIN DBMS_SCHEDULER.RUN_JOB(JOB_SYNC_ORDER, FALSE); END; /RUN_JOB的第二个参数设为FALSE表示立即执行不等待。执行后查日志SELECT * FROM MIGRATION_LOG WHERE TASK_NAME ORDER_SYNC ORDER BY LOG_TIME DESC LIMIT 5;确认执行成功后启用任务BEGIN DBMS_SCHEDULER.ENABLE(JOB_SYNC_ORDER); END; /5.4 日常运维中需要关注的几个指标任务跑起来之后日常运维需要关注几个关键指标。迁移延迟从源数据产生到同步到达梦的时间差。如果延迟持续增大说明迁移速度跟不上数据产生速度需要调整批次大小或者调度频率。单次迁移行数如果某次迁移的行数突然暴增可能是源库有批量操作需要确认是否正常。如果行数持续为0可能是水位线没有正确更新或者源表没有新数据。执行时长如果执行时长突然变长可能是源表数据量增长、索引失效、或者锁等待。需要结合数据库的等待事件来分析。错误日志定期检查MIGRATION_LOG表中STATUS FAILED的记录及时处理。-- 查询最近7天的迁移统计 SELECT TRUNC(LOG_TIME) AS LOG_DATE, COUNT(*) AS RUN_COUNT, SUM(CASE WHEN STATUS SUCCESS THEN 1 ELSE 0 END) AS SUCCESS_COUNT, SUM(CASE WHEN STATUS FAILED THEN 1 ELSE 0 END) AS FAILED_COUNT, SUM(ROW_COUNT) AS TOTAL_ROWS, AVG(DURATION) AS AVG_DURATION_SEC FROM MIGRATION_LOG WHERE TASK_NAME ORDER_SYNC AND LOG_TIME SYSDATE - 7 GROUP BY TRUNC(LOG_TIME) ORDER BY LOG_DATE DESC;这个查询可以直观地看到每天的迁移次数、成功失败情况、总行数和平均耗时。如果发现某天失败次数较多可以进一步查当天的错误信息。6. 一些值得考虑的优化方向6.1 并行迁移让多个任务同时干活如果单次迁移的数据量很大单线程执行时间太长可以考虑把迁移任务拆成多个并行执行的子任务。比如按订单ID的哈希值取模分成4个子任务每个子任务负责一部分数据同时执行。达梦的DBMS_SCHEDULER支持创建多个Job只要它们的执行时间不冲突就可以并行运行。但并行迁移需要注意几个问题一是水位线的管理要按子任务分开不能共用一个水位线二是并行插入可能加剧目标表的锁竞争需要评估目标表的写入能力三是并行任务的错误处理要独立一个子任务失败不应该影响其他子任务。6.2 迁移失败后的自动重试存储过程执行失败的原因有很多有些是暂时性的比如锁等待超时、网络抖动有些是永久性的比如表不存在、字段类型不匹配。对于暂时性的失败可以配置自动重试。达梦的DBMS_SCHEDULER本身不直接支持失败重试但可以在存储过程内部实现重试逻辑V_RETRY_COUNT INT : 0; V_MAX_RETRY INT : 3; RETRY_LOOP LOOP BEGIN -- 迁移逻辑 ... EXIT RETRY_LOOP; EXCEPTION WHEN OTHERS THEN V_RETRY_COUNT : V_RETRY_COUNT 1; IF V_RETRY_COUNT V_MAX_RETRY THEN RAISE; END IF; -- 等待5秒后重试 DBMS_LOCK.SLEEP(5); END; END LOOP RETRY_LOOP;DBMS_LOCK.SLEEP是达梦提供的休眠函数参数是秒数。重试之前先休眠几秒给数据库一个缓冲的时间避免立即重试又撞上同样的锁。6.3 迁移数据的校验与对账数据迁移完成之后怎么确认迁移的数据是对的最直接的方式是对账——比较源表和目标表在某个时间范围内的记录数和关键字段的汇总值。-- 源表统计 SELECT COUNT(*) AS CNT, SUM(AMOUNT) AS TOTAL_AMOUNT FROM SRC_ORDER WHERE UPDATE_TIME BETWEEN 2024-01-01 AND 2024-01-02; -- 目标表统计 SELECT COUNT(*) AS CNT, SUM(AMOUNT) AS TOTAL_AMOUNT FROM TGT_ORDER WHERE UPDATE_TIME BETWEEN 2024-01-01 AND 2024-01-02;如果两边对不上就需要进一步排查。常见的对不上的原因包括源表在迁移过程中有新数据写入、目标表的唯一约束导致部分数据被跳过、类型转换导致精度丢失等。对账逻辑也可以封装成存储过程每天迁移完成后自动执行发现不一致就发告警。这样就不需要人工去检查了。6.4 历史数据的归档清理迁移任务跑久了目标表的数据会越来越多。如果目标表只是用来做报表分析历史数据的查询频率很低可以考虑定期归档。比如把一年前的数据从目标表转移到历史表目标表只保留最近一年的数据。归档操作也可以做成存储过程用DBMS_SCHEDULER调度。归档的时候要注意先插入历史表确认插入成功后再删除目标表的数据整个过程放在一个事务里。如果历史表的数据量也很大同样需要分批提交。CREATE OR REPLACE PROCEDURE PROC_ARCHIVE_ORDER( P_KEEP_MONTHS IN INT DEFAULT 12 ) AS V_CUTOFF_DATE TIMESTAMP; BEGIN V_CUTOFF_DATE : ADD_MONTHS(SYSDATE, -P_KEEP_MONTHS); -- 插入历史表 INSERT INTO TGT_ORDER_HIS SELECT * FROM TGT_ORDER WHERE UPDATE_TIME V_CUTOFF_DATE; -- 删除目标表数据 DELETE FROM TGT_ORDER WHERE UPDATE_TIME V_CUTOFF_DATE; COMMIT; END; /归档和迁移最好不要在同一个时间窗口执行避免资源竞争。可以把归档安排在迁移完成之后比如凌晨4点。7. 写在最后的一些个人体会这套方案我在几个项目里实际用过整体来说是稳定的但也不是没有代价。最大的代价是维护成本——存储过程不像应用程序代码那样有版本管理、单元测试、CI/CD它的变更和回滚都更原始。所以我的建议是存储过程的代码一定要纳入版本管理每次变更都要有记录变更前先在测试环境验证。另一个体会是日志和监控比迁移逻辑本身更重要。迁移逻辑写错了可以改但如果迁移失败了没人知道那问题就大了。所以我在每个项目里都会花不少时间在日志表和监控查询上确保任何异常都能被及时发现。还有一个容易被忽略的点是水位线的初始值。如果水位线设置得太早第一次迁移会拉取大量历史数据可能导致执行时间过长如果设置得太晚又会漏掉一部分数据。我的做法是先用源表的MIN(UPDATE_TIME)作为初始水位线然后根据数据量决定是否需要先做一次历史数据的一次性迁移。最后说一个实际踩过的坑达梦的MERGE语句在某些版本下如果USING子查询返回了重复的记录会报错。所以在用MERGE之前一定要确保USING子查询里的数据在ON条件的字段上是唯一的。如果源表可能存在重复先用GROUP BY或者ROW_NUMBER()去重。这个坑我在一个项目里踩过排查了大半天才发现是源表有重复数据导致的。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进