news 2026/9/19 1:23:24

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
达梦数据库存储过程与定时任务实现数据自动迁移方案

数据迁移这件事,做过一次的人都知道,最怕的不是迁移本身,而是迁移完之后业务方隔三差五来找你:“昨天的数据怎么还没同步过来?”你打开工具手动跑一遍,数据是过来了,但明天呢?后天呢?这种重复劳动做多了,人就会开始琢磨——能不能让数据库自己把这件事干了?

达梦数据库作为国产数据库中的主力选手,在越来越多的生产环境中承担着核心业务数据的存储和管理。但很多团队在从其他数据库迁移到达梦之后,往往只关注“数据能不能导进去”,而忽略了“数据能不能持续自动地流转”。手动导数在一次性迁移场景下没问题,可一旦涉及周期性同步、跨库数据汇总、历史数据归档这类需求,靠人工操作就完全不现实了。这篇文章要聊的,就是怎么利用达梦的存储过程配合定时任务,搭建一套能自己跑起来的数据自动迁移方案。

这套方案适合谁?如果你正在使用达梦数据库,手头有周期性数据同步的需求,又不想引入额外的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 => 'FREQ=DAILY; BYHOUR=2; BYMINUTE=0; BYSECOND=0', ENABLED => TRUE, COMMENTS => '订单数据每日凌晨2点自动迁移' ); END; /

几个关键参数需要解释一下。

JOB_TYPE指定任务类型,这里用的是STORED_PROCEDURE,表示执行一个存储过程。达梦还支持PLSQL_BLOCK,可以直接写一段PL/SQL代码块,适合逻辑比较简单的场景。如果迁移逻辑不复杂,用PLSQL_BLOCK可以省去创建存储过程的步骤。

REPEAT_INTERVAL是调度频率,用的是日历表达式语法。FREQ=DAILY; BYHOUR=2; BYMINUTE=0; BYSECOND=0表示每天凌晨2点整执行。这个表达式可以组合出很复杂的调度策略,比如FREQ=WEEKLY; BYDAY=MON,TUE,WED,THU,FRI; BYHOUR=3表示周一到周五每天凌晨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_DATEREPEAT_INTERVAL是配合使用的。如果START_DATE设的是SYSDATE,而当前时间是下午3点,REPEAT_INTERVAL设的是FREQ=DAILY; BYHOUR=2,那么第一次执行会在明天凌晨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;

另外,如果存储过程内部有INSERTUPDATEDELETE操作,还需要确保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里SQLCODESQLERRM在某些异常场景下返回的值和达梦8不一致,导致错误日志记录不准确。解决办法是在存储过程里用WHEN OTHERS THEN捕获异常后,同时记录SQLCODE和自定义的错误描述,而不是完全依赖SQLERRM

另一个版本差异是调度任务的时区处理。达梦8的某些版本在START_DATE的处理上会考虑数据库的时区设置,而达梦7则统一按服务器本地时间处理。如果从达梦7升级到达梦8,原有的调度任务可能需要重新调整时间配置。

建议:在正式部署之前,先在测试环境用目标版本做一次完整的验证,包括存储过程编译、Job创建、手动执行、自动调度等环节。不要假设在开发环境能跑通的代码在生产环境也一定能跑通。

5. 从手动到自动:一个完整的落地案例

5.1 场景描述与表结构设计

假设有一个订单系统,源库是MySQL,目标库是达梦。业务要求每天凌晨把MySQL中当天的订单数据同步到达梦的分析库中,供报表系统查询。

源表结构(MySQL):

CREATE 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 => 'FREQ=DAILY; BYHOUR=2; BYMINUTE=0; BYSECOND=0', 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()去重。这个坑我在一个项目里踩过,排查了大半天才发现是源表有重复数据导致的。

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

郑州A.O.史密斯热水器故障维修电话|内胆漏水上门排查|欧米到家报修热线

洗澡时热水忽冷忽热、燃气热水器打不着火、电热水器加热慢、空气能热水不够用、太阳能控制器报警……这些问题表面上都指向“没有热水”&#xff0c;实际背后却可能涉及水路、电路、燃气、燃烧、排烟、温控、传感器、安装环境及长期维护等多个环节。真正专业的热水器维修&#…

作者头像 李华