You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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流水线的方式

  1. 在ADF中添加一个Oracle数据库链接,连接到你的Oracle实例;
  2. 使用存储过程活动调用上述PL/SQL存储过程;
  3. 如果存储过程返回游标,直接用ADF的复制活动读取游标结果集,同步到目标数据源;
  4. 将存储过程执行和数据复制步骤编排成流水线,设置每日调度执行。

注意事项

  • 对于大数据量的表,建议给时间戳字段或日志表的CHANGE_TIMESTAMP字段建索引,提升查询性能;
  • 如果需要处理删除操作,方案二的日志表会记录删除动作,ADF可以根据OPERATION_TYPE在目标库执行对应操作;
  • 同步控制表要做好权限管理,避免误修改同步时间导致数据重复或丢失。

内容的提问来源于stack exchange,提问作者RK.

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.07 23:55:04