Oracle EDW CDC替代方案咨询:PL/SQL实现每日增量数据获取
可以用PL/SQL实现每日增量数据捕获,适配ADF流水线集成
完全可以通过PL/SQL存储过程实现每日增量数据获取,下面是两种实用方案,适配不同场景:
方案一:基于时间戳字段的增量提取(推荐优先使用)
如果你的业务表已经包含LAST_UPDATED_DATE、CREATE_DATE这类记录数据变更时间的字段,直接用这个方案最省心:
1. 先建一个同步控制表
用来记录每个表的上次同步时间,避免重复抓取:
CREATE TABLE SYNC_CONTROL ( TABLE_NAME VARCHAR2(100) PRIMARY KEY, LAST_SYNC_TIMESTAMP TIMESTAMP DEFAULT SYSTIMESTAMP ); -- 初始化目标表的同步时间 INSERT INTO SYNC_CONTROL (TABLE_NAME) VALUES ('YOUR_BUSINESS_TABLE'); COMMIT;
2. 编写PL/SQL存储过程
这个过程会读取上次同步时间,抓取增量数据,然后更新同步时间:
CREATE OR REPLACE PROCEDURE GET_INCREMENTAL_DATA(p_table_name IN VARCHAR2, p_output_cursor OUT SYS_REFCURSOR) IS v_last_sync TIMESTAMP; BEGIN -- 获取上次同步时间 SELECT LAST_SYNC_TIMESTAMP INTO v_last_sync FROM SYNC_CONTROL WHERE TABLE_NAME = p_table_name; -- 抓取增量数据(这里以SELECT为例,也可以插入到临时表供ADF读取) OPEN p_output_cursor FOR SELECT * FROM YOUR_BUSINESS_TABLE WHERE LAST_UPDATED_DATE > v_last_sync; -- 更新同步时间为当前时间 UPDATE SYNC_CONTROL SET LAST_SYNC_TIMESTAMP = SYSTIMESTAMP WHERE TABLE_NAME = p_table_name; COMMIT; EXCEPTION WHEN NO_DATA_FOUND THEN RAISE_APPLICATION_ERROR(-20001, '表' || p_table_name || '未在同步控制表中初始化'); WHEN OTHERS THEN ROLLBACK; RAISE; END; /
方案二:基于触发器+变更日志表的捕获(无时间戳字段时使用)
如果业务表没有时间戳字段,需要通过触发器记录所有变更操作:
1. 创建变更日志表
CREATE TABLE CHANGE_LOG ( LOG_ID NUMBER GENERATED ALWAYS AS IDENTITY PRIMARY KEY, TABLE_NAME VARCHAR2(100), OPERATION_TYPE VARCHAR2(10), -- INSERT/UPDATE/DELETE CHANGE_TIMESTAMP TIMESTAMP DEFAULT SYSTIMESTAMP, OLD_DATA CLOB, -- 存储旧数据(可选) NEW_DATA CLOB -- 存储新数据 );
2. 给业务表建触发器
以INSERT/UPDATE/DELETE为例:
CREATE OR REPLACE TRIGGER YOUR_TABLE_TRIGGER AFTER INSERT OR UPDATE OR DELETE ON YOUR_BUSINESS_TABLE FOR EACH ROW DECLARE v_operation VARCHAR2(10); v_old_data CLOB; v_new_data CLOB; BEGIN -- 判断操作类型 IF INSERTING THEN v_operation := 'INSERT'; v_new_data := JSON_OBJECT_T(:NEW).TO_CLOB(); ELSIF UPDATING THEN v_operation := 'UPDATE'; v_old_data := JSON_OBJECT_T(:OLD).TO_CLOB(); v_new_data := JSON_OBJECT_T(:NEW).TO_CLOB(); ELSIF DELETING THEN v_operation := 'DELETE'; v_old_data := JSON_OBJECT_T(:OLD).TO_CLOB(); END IF; -- 插入日志 INSERT INTO CHANGE_LOG (TABLE_NAME, OPERATION_TYPE, OLD_DATA, NEW_DATA) VALUES ('YOUR_BUSINESS_TABLE', v_operation, v_old_data, v_new_data); END; /
3. 编写提取增量的存储过程
CREATE OR REPLACE PROCEDURE GET_DAILY_CHANGES(p_output_cursor OUT SYS_REFCURSOR) IS v_yesterday_start TIMESTAMP; v_today_start TIMESTAMP; BEGIN -- 计算昨日0点到今日0点的时间范围 v_yesterday_start := TRUNC(SYSTIMESTAMP - INTERVAL '1' DAY); v_today_start := TRUNC(SYSTIMESTAMP); OPEN p_output_cursor FOR SELECT * FROM CHANGE_LOG WHERE CHANGE_TIMESTAMP BETWEEN v_yesterday_start AND v_today_start ORDER BY CHANGE_TIMESTAMP; END; /
集成到ADF流水线的方式
- 在ADF中添加一个Oracle数据库链接,连接到你的Oracle实例;
- 使用存储过程活动调用上述PL/SQL存储过程;
- 如果存储过程返回游标,直接用ADF的复制活动读取游标结果集,同步到目标数据源;
- 将存储过程执行和数据复制步骤编排成流水线,设置每日调度执行。
注意事项
- 对于大数据量的表,建议给时间戳字段或日志表的
CHANGE_TIMESTAMP字段建索引,提升查询性能; - 如果需要处理删除操作,方案二的日志表会记录删除动作,ADF可以根据
OPERATION_TYPE在目标库执行对应操作; - 同步控制表要做好权限管理,避免误修改同步时间导致数据重复或丢失。
内容的提问来源于stack exchange,提问作者RK.
相关产品推荐
相关产品推荐

