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

如何并行执行存储过程中游标记录的批量处理任务

大规模关联表并行迁移的存储过程实现方案

表结构定义

首先给出两张关联表的创建语句:

CREATE TABLE address (adr_id, ver_id, address) AS
SELECT 1, 1, 'newYork' FROM DUAL UNION ALL
SELECT 1, 2, 'newYork' FROM DUAL UNION ALL
SELECT 1, 3, 'newYork' FROM DUAL UNION ALL
SELECT 4, 1, 'Washington' FROM DUAL UNION ALL
SELECT 4, 2, 'Washington' FROM DUAL;

CREATE TABLE employee (emp_id,adr_id,ver_id) AS
SELECT 100,1, 1 FROM DUAL UNION ALL
SELECT 200,1, 2 FROM DUAL UNION ALL
SELECT 300,1, 3 FROM DUAL UNION ALL
SELECT 400,4, 1 FROM DUAL UNION ALL
SELECT 500,4, 2 FROM DUAL;

说明:两张表均有数十亿条记录,外键约束已放宽,需通过并行执行提升处理效率。

核心任务要求

存储过程需要完成以下操作:

  • 筛选出地址为newYork的所有独立adr_id分组,存入游标
  • 遍历游标中的每个分组,执行:
    1. 查询该adr_id下所有address表记录(即同地址的不同版本)
    2. 插入一条ver_id=0的新地址记录,地址与原分组一致,并获取新的adr_id
    3. 将employee表中关联原adr_id的所有记录,更新为新生成的adr_id
    4. 删除原adr_id对应的所有旧地址记录

并行处理实现方案

针对数十亿级数据,单进程处理效率极低,推荐使用Oracle的DBMS_PARALLEL_EXECUTE包实现分块并行处理——它能将大任务拆分为多个独立小任务,分配到不同并行进程执行,避免锁冲突和性能瓶颈。

具体实现步骤

1. 创建并行任务与数据分块

先创建并行任务,按adr_id对需要处理的数据进行分块,确保每个分块包含独立的adr_id分组,避免并行操作时的冲突:

DECLARE
  l_task_name VARCHAR2(100) := 'ADDR_MIGRATION_TASK';
  l_sql_stmt  VARCHAR2(1000);
BEGIN
  -- 初始化并行任务
  DBMS_PARALLEL_EXECUTE.CREATE_TASK(task_name => l_task_name);

  -- 按adr_id分组生成分块,此处仅针对newYork地址
  l_sql_stmt := 'SELECT adr_id, adr_id FROM address WHERE address = ''newYork'' GROUP BY adr_id';
  DBMS_PARALLEL_EXECUTE.CREATE_CHUNKS_BY_SQL(
    task_name => l_task_name,
    sql_stmt  => l_sql_stmt,
    by_rowid  => FALSE
  );
END;
/

2. 编写分块处理的存储过程

创建一个带自治事务的处理过程,每个并行进程独立处理一个分块,避免长事务阻塞:

CREATE OR REPLACE PROCEDURE process_address_chunk(p_start_id IN NUMBER, p_end_id IN NUMBER) IS
  l_new_adr_id NUMBER;
  l_address    VARCHAR2(100);
BEGIN
  -- 遍历分块内的每个adr_id分组
  FOR rec IN (SELECT adr_id, address FROM address WHERE adr_id BETWEEN p_start_id AND p_end_id GROUP BY adr_id, address) LOOP
    l_address := rec.address;

    -- 插入新的ver_id=0地址记录,用序列生成新adr_id
    INSERT INTO address(adr_id, ver_id, address)
    VALUES(adr_id_seq.NEXTVAL, 0, l_address)
    RETURNING adr_id INTO l_new_adr_id;

    -- 更新关联的employee记录
    UPDATE employee
    SET adr_id = l_new_adr_id
    WHERE adr_id = rec.adr_id;

    -- 删除原adr_id的所有旧记录
    DELETE FROM address
    WHERE adr_id = rec.adr_id;

    -- 独立提交,释放锁资源
    COMMIT;
  END LOOP;
EXCEPTION
  WHEN OTHERS THEN
    DBMS_OUTPUT.PUT_LINE('处理adr_id范围 ' || p_start_id || '-' || p_end_id || ' 失败: ' || SQLERRM);
    ROLLBACK;
END;
/

注意:需提前创建adr_id_seq序列用于生成新的adr_id值:

CREATE SEQUENCE adr_id_seq START WITH 11 INCREMENT BY 1;

3. 启动并行任务并清理

调用RUN_TASK启动并行执行,并行度根据服务器CPU核心数调整(比如8核服务器设为8):

DECLARE
  l_task_name VARCHAR2(100) := 'ADDR_MIGRATION_TASK';
BEGIN
  DBMS_PARALLEL_EXECUTE.RUN_TASK(
    task_name      => l_task_name,
    sql_stmt       => 'BEGIN process_address_chunk(:start_id, :end_id); END;',
    language_flag  => DBMS_SQL.NATIVE,
    parallel_level => 8 -- 并行度按需调整
  );

  -- 检查任务状态
  IF DBMS_PARALLEL_EXECUTE.TASK_STATUS(l_task_name) = DBMS_PARALLEL_EXECUTE.FINISHED THEN
    DBMS_OUTPUT.PUT_LINE('并行任务执行完成');
  ELSE
    DBMS_OUTPUT.PUT_LINE('并行任务执行失败,请检查错误日志');
  END;

  -- 清理任务资源
  DBMS_PARALLEL_EXECUTE.DROP_TASK(task_name => l_task_name);
END;
/

关键优化建议

  • 分块策略:如果数据分布不均,可改用哈希分块(CREATE_CHUNKS_BY_HASH),确保每个分块数据量均匀。
  • 索引优化:在address(adr_id, address)和employee(adr_id)上创建索引,大幅加速查询、更新和删除操作。
  • 并行度控制:并行度不要超过CPU核心数的1.5倍,避免IO或CPU资源耗尽。
  • 增量处理:如果数据持续写入,可按时间范围分批次处理,避免锁表影响业务。

预期执行结果

Address表

adr_idver_idaddress
110newYork
120Washington

Employee表

emp_idadr_idver_id
100111
200112
300113
400121
500122

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 23:01:04